70 changed files with 1938 additions and 252 deletions
@ -0,0 +1,949 @@ |
|||
/** |
|||
* Copyright © 2016-2024 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.actors.rule; |
|||
|
|||
import org.junit.jupiter.api.BeforeEach; |
|||
import org.junit.jupiter.api.Test; |
|||
import org.junit.jupiter.api.extension.ExtendWith; |
|||
import org.junit.jupiter.params.ParameterizedTest; |
|||
import org.junit.jupiter.params.provider.Arguments; |
|||
import org.junit.jupiter.params.provider.MethodSource; |
|||
import org.junit.jupiter.params.provider.ValueSource; |
|||
import org.mockito.ArgumentCaptor; |
|||
import org.mockito.Mock; |
|||
import org.mockito.Mockito; |
|||
import org.mockito.junit.jupiter.MockitoExtension; |
|||
import org.thingsboard.server.actors.ActorSystemContext; |
|||
import org.thingsboard.server.actors.TbActorRef; |
|||
import org.thingsboard.server.actors.ruleChain.DefaultTbContext; |
|||
import org.thingsboard.server.actors.ruleChain.RuleChainOutputMsg; |
|||
import org.thingsboard.server.actors.ruleChain.RuleNodeCtx; |
|||
import org.thingsboard.server.actors.ruleChain.RuleNodeToRuleChainTellNextMsg; |
|||
import org.thingsboard.server.cluster.TbClusterService; |
|||
import org.thingsboard.server.common.data.DataConstants; |
|||
import org.thingsboard.server.common.data.debug.DebugSettings; |
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
import org.thingsboard.server.common.data.id.RuleChainId; |
|||
import org.thingsboard.server.common.data.id.RuleNodeId; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.data.msg.TbMsgType; |
|||
import org.thingsboard.server.common.data.msg.TbNodeConnectionType; |
|||
import org.thingsboard.server.common.data.rule.RuleNode; |
|||
import org.thingsboard.server.common.msg.TbMsg; |
|||
import org.thingsboard.server.common.msg.TbMsgMetaData; |
|||
import org.thingsboard.server.common.msg.TbMsgProcessingStackItem; |
|||
import org.thingsboard.server.common.msg.queue.ServiceType; |
|||
import org.thingsboard.server.common.msg.queue.TbMsgCallback; |
|||
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; |
|||
import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineMsg; |
|||
import org.thingsboard.server.queue.common.SimpleTbQueueCallback; |
|||
|
|||
import java.util.Collections; |
|||
import java.util.List; |
|||
import java.util.Set; |
|||
import java.util.UUID; |
|||
import java.util.concurrent.TimeUnit; |
|||
import java.util.function.Consumer; |
|||
import java.util.stream.Collectors; |
|||
import java.util.stream.Stream; |
|||
|
|||
import static org.assertj.core.api.Assertions.assertThat; |
|||
import static org.mockito.ArgumentMatchers.any; |
|||
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; |
|||
import static org.mockito.Mockito.never; |
|||
import static org.mockito.Mockito.times; |
|||
|
|||
@SuppressWarnings("ResultOfMethodCallIgnored") |
|||
@ExtendWith(MockitoExtension.class) |
|||
class DefaultTbContextTest { |
|||
|
|||
private final String EXCEPTION_MSG = "Some runtime exception!"; |
|||
private final RuntimeException EXCEPTION = new RuntimeException(EXCEPTION_MSG); |
|||
|
|||
private final TenantId TENANT_ID = TenantId.fromUUID(UUID.fromString("c7bf4c85-923c-4688-a4b5-0f8a0feb7cd5")); |
|||
private final RuleNodeId RULE_NODE_ID = new RuleNodeId(UUID.fromString("1ca5e2ef-1309-41d9-bafa-709e9df0e2a6")); |
|||
private final RuleChainId RULE_CHAIN_ID = new RuleChainId(UUID.fromString("b87c4123-f9f2-41a6-9a09-e3a5b6580b11")); |
|||
|
|||
@Mock |
|||
private ActorSystemContext mainCtxMock; |
|||
@Mock |
|||
private RuleNodeCtx nodeCtxMock; |
|||
@Mock |
|||
private TbActorRef chainActorMock; |
|||
|
|||
private DefaultTbContext defaultTbContext; |
|||
|
|||
@BeforeEach |
|||
public void setUp() { |
|||
defaultTbContext = new DefaultTbContext(mainCtxMock, "Test rule chain name", nodeCtxMock); |
|||
} |
|||
|
|||
@Test |
|||
public void givenDebugFailuresEvents_whenTellSuccess_thenVerifyDebugOutputNotPersisted() { |
|||
// GIVEN
|
|||
var callbackMock = mock(TbMsgCallback.class); |
|||
var msg = getTbMsgWithCallback(callbackMock); |
|||
var ruleNode = new RuleNode(RULE_NODE_ID); |
|||
ruleNode.setRuleChainId(RULE_CHAIN_ID); |
|||
ruleNode.setDebugSettings(DebugSettings.failures()); |
|||
given(nodeCtxMock.getSelf()).willReturn(ruleNode); |
|||
given(nodeCtxMock.getChainActor()).willReturn(chainActorMock); |
|||
|
|||
// WHEN
|
|||
defaultTbContext.tellSuccess(msg); |
|||
|
|||
// THEN
|
|||
then(nodeCtxMock).should().getChainActor(); |
|||
then(nodeCtxMock).shouldHaveNoMoreInteractions(); |
|||
then(mainCtxMock).shouldHaveNoInteractions(); |
|||
checkTellNextCommonLogic(callbackMock, TbNodeConnectionType.SUCCESS, msg); |
|||
} |
|||
|
|||
@Test |
|||
public void givenDebugFailuresEventsAndSuccessConnection_whenTellNext_thenVerifyDebugOutputNotPersisted() { |
|||
// GIVEN
|
|||
var callbackMock = mock(TbMsgCallback.class); |
|||
var msg = getTbMsgWithCallback(callbackMock); |
|||
var ruleNode = new RuleNode(RULE_NODE_ID); |
|||
ruleNode.setRuleChainId(RULE_CHAIN_ID); |
|||
ruleNode.setDebugSettings(DebugSettings.failures()); |
|||
given(nodeCtxMock.getSelf()).willReturn(ruleNode); |
|||
given(nodeCtxMock.getChainActor()).willReturn(chainActorMock); |
|||
|
|||
// WHEN
|
|||
defaultTbContext.tellNext(msg, TbNodeConnectionType.SUCCESS); |
|||
|
|||
// THEN
|
|||
then(nodeCtxMock).should().getChainActor(); |
|||
then(nodeCtxMock).shouldHaveNoMoreInteractions(); |
|||
then(mainCtxMock).shouldHaveNoInteractions(); |
|||
checkTellNextCommonLogic(callbackMock, TbNodeConnectionType.SUCCESS, msg); |
|||
} |
|||
|
|||
@MethodSource |
|||
@ParameterizedTest |
|||
void givenDebugFailuresEventsAndConnections_whenTellNext_thenVerifyDebugOutputPersisted(Set<String> connections) { |
|||
// GIVEN
|
|||
var callbackMock = mock(TbMsgCallback.class); |
|||
var msg = getTbMsgWithCallback(callbackMock); |
|||
var ruleNode = new RuleNode(RULE_NODE_ID); |
|||
ruleNode.setRuleChainId(RULE_CHAIN_ID); |
|||
ruleNode.setDebugSettings(DebugSettings.failures()); |
|||
given(nodeCtxMock.getTenantId()).willReturn(TENANT_ID); |
|||
given(nodeCtxMock.getSelf()).willReturn(ruleNode); |
|||
given(nodeCtxMock.getChainActor()).willReturn(chainActorMock); |
|||
|
|||
// WHEN
|
|||
defaultTbContext.tellNext(msg, connections); |
|||
|
|||
// THEN
|
|||
then(nodeCtxMock).should().getChainActor(); |
|||
then(nodeCtxMock).shouldHaveNoMoreInteractions(); |
|||
then(mainCtxMock).should().persistDebugOutput(TENANT_ID, RULE_NODE_ID, msg, TbNodeConnectionType.FAILURE, null, null); |
|||
then(mainCtxMock).shouldHaveNoMoreInteractions(); |
|||
checkTellNextCommonLogic(callbackMock, connections, msg); |
|||
} |
|||
|
|||
private static Stream<Set<String>> givenDebugFailuresEventsAndConnections_whenTellNext_thenVerifyDebugOutputPersisted() { |
|||
return Stream.of( |
|||
Collections.singleton(TbNodeConnectionType.FAILURE), |
|||
Set.of(TbNodeConnectionType.FAILURE, TbNodeConnectionType.SUCCESS) |
|||
); |
|||
} |
|||
|
|||
@MethodSource |
|||
@ParameterizedTest |
|||
void givenDebugDisabledAndConnections_whenTellNext_thenVerifyDebugOutputNotPersisted(Set<String> connections) { |
|||
// GIVEN
|
|||
var callbackMock = mock(TbMsgCallback.class); |
|||
var msg = getTbMsgWithCallback(callbackMock); |
|||
var ruleNode = new RuleNode(RULE_NODE_ID); |
|||
ruleNode.setRuleChainId(RULE_CHAIN_ID); |
|||
ruleNode.setDebugSettings(DebugSettings.off()); |
|||
given(nodeCtxMock.getSelf()).willReturn(ruleNode); |
|||
given(nodeCtxMock.getChainActor()).willReturn(chainActorMock); |
|||
|
|||
// WHEN
|
|||
defaultTbContext.tellNext(msg, connections); |
|||
|
|||
// THEN
|
|||
then(nodeCtxMock).should().getChainActor(); |
|||
then(nodeCtxMock).shouldHaveNoMoreInteractions(); |
|||
then(mainCtxMock).shouldHaveNoInteractions(); |
|||
checkTellNextCommonLogic(callbackMock, connections, msg); |
|||
} |
|||
|
|||
private static Stream<Set<String>> givenDebugDisabledAndConnections_whenTellNext_thenVerifyDebugOutputNotPersisted() { |
|||
return Stream.of( |
|||
Collections.singleton(TbNodeConnectionType.FAILURE), |
|||
Collections.singleton(TbNodeConnectionType.SUCCESS), |
|||
Set.of(TbNodeConnectionType.FAILURE, TbNodeConnectionType.SUCCESS) |
|||
); |
|||
} |
|||
|
|||
@MethodSource |
|||
@ParameterizedTest |
|||
void givenDebugAllEventsAndConnection_whenTellNext_thenVerifyDebugOutputPersisted(String connection) { |
|||
// GIVEN
|
|||
var callbackMock = mock(TbMsgCallback.class); |
|||
var msg = getTbMsgWithCallback(callbackMock); |
|||
var ruleNode = new RuleNode(RULE_NODE_ID); |
|||
ruleNode.setRuleChainId(RULE_CHAIN_ID); |
|||
ruleNode.setDebugSettings(DebugSettings.until(getUntilTime())); |
|||
given(nodeCtxMock.getTenantId()).willReturn(TENANT_ID); |
|||
given(nodeCtxMock.getSelf()).willReturn(ruleNode); |
|||
given(nodeCtxMock.getChainActor()).willReturn(chainActorMock); |
|||
|
|||
// WHEN
|
|||
defaultTbContext.tellNext(msg, connection); |
|||
|
|||
// THEN
|
|||
then(nodeCtxMock).should().getChainActor(); |
|||
then(nodeCtxMock).shouldHaveNoMoreInteractions(); |
|||
then(mainCtxMock).should().persistDebugOutput(TENANT_ID, RULE_NODE_ID, msg, connection, null, null); |
|||
then(mainCtxMock).shouldHaveNoMoreInteractions(); |
|||
checkTellNextCommonLogic(callbackMock, connection, msg); |
|||
} |
|||
|
|||
private static Stream<String> givenDebugAllEventsAndConnection_whenTellNext_thenVerifyDebugOutputPersisted() { |
|||
return failureAndSuccessConnection(); |
|||
} |
|||
|
|||
@Test |
|||
public void givenDebugAllEventsAndFailureAndSuccessConnection_whenTellNext_thenVerifyDebugOutputPersistedForAllEvents() { |
|||
// GIVEN
|
|||
var callbackMock = mock(TbMsgCallback.class); |
|||
var msg = getTbMsgWithCallback(callbackMock); |
|||
var ruleNode = new RuleNode(RULE_NODE_ID); |
|||
ruleNode.setRuleChainId(RULE_CHAIN_ID); |
|||
ruleNode.setDebugSettings(DebugSettings.until(getUntilTime())); |
|||
given(nodeCtxMock.getTenantId()).willReturn(TENANT_ID); |
|||
given(nodeCtxMock.getSelf()).willReturn(ruleNode); |
|||
given(nodeCtxMock.getChainActor()).willReturn(chainActorMock); |
|||
|
|||
// WHEN
|
|||
Set<String> connections = failureAndSuccessConnection().collect(Collectors.toSet()); |
|||
defaultTbContext.tellNext(msg, connections); |
|||
|
|||
// THEN
|
|||
then(nodeCtxMock).should().getChainActor(); |
|||
then(nodeCtxMock).shouldHaveNoMoreInteractions(); |
|||
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(), nullable(Throwable.class), nullable(String.class)); |
|||
then(mainCtxMock).shouldHaveNoMoreInteractions(); |
|||
assertThat(nodeConnectionsCaptor.getAllValues()).hasSize(wantedNumberOfInvocations); |
|||
assertThat(nodeConnectionsCaptor.getAllValues()).containsExactlyInAnyOrderElementsOf(connections); |
|||
checkTellNextCommonLogic(callbackMock, connections, msg); |
|||
} |
|||
|
|||
@MethodSource |
|||
@ParameterizedTest |
|||
void givenDebugAllThenOnlyFailureEventsAndConnection_whenTellNext_thenVerifyDebugOutputPersisted(String connection) { |
|||
// GIVEN
|
|||
var callbackMock = mock(TbMsgCallback.class); |
|||
var msg = getTbMsgWithCallback(callbackMock); |
|||
var ruleNode = new RuleNode(RULE_NODE_ID); |
|||
ruleNode.setRuleChainId(RULE_CHAIN_ID); |
|||
ruleNode.setDebugSettings(DebugSettings.until(getUntilTime())); |
|||
given(nodeCtxMock.getTenantId()).willReturn(TENANT_ID); |
|||
given(nodeCtxMock.getSelf()).willReturn(ruleNode); |
|||
given(nodeCtxMock.getChainActor()).willReturn(chainActorMock); |
|||
|
|||
// WHEN
|
|||
defaultTbContext.tellNext(msg, connection); |
|||
|
|||
// THEN
|
|||
then(nodeCtxMock).should().getChainActor(); |
|||
then(nodeCtxMock).shouldHaveNoMoreInteractions(); |
|||
then(mainCtxMock).should().persistDebugOutput(TENANT_ID, RULE_NODE_ID, msg, connection, null, null); |
|||
then(mainCtxMock).shouldHaveNoMoreInteractions(); |
|||
checkTellNextCommonLogic(callbackMock, connection, msg); |
|||
} |
|||
|
|||
private static Stream<String> givenDebugAllThenOnlyFailureEventsAndConnection_whenTellNext_thenVerifyDebugOutputPersisted() { |
|||
return failureAndSuccessConnection(); |
|||
} |
|||
|
|||
@Test |
|||
public void givenDebugAllThenOnlyEventsAndFailureAndSuccessConnection_whenTellNext_thenVerifyDebugOutputPersistedForAllEvents() { |
|||
// GIVEN
|
|||
var callbackMock = mock(TbMsgCallback.class); |
|||
var msg = getTbMsgWithCallback(callbackMock); |
|||
var ruleNode = new RuleNode(RULE_NODE_ID); |
|||
ruleNode.setRuleChainId(RULE_CHAIN_ID); |
|||
ruleNode.setDebugSettings(DebugSettings.failuresOrUntil(getUntilTime())); |
|||
given(nodeCtxMock.getTenantId()).willReturn(TENANT_ID); |
|||
given(nodeCtxMock.getSelf()).willReturn(ruleNode); |
|||
given(nodeCtxMock.getChainActor()).willReturn(chainActorMock); |
|||
|
|||
// WHEN
|
|||
Set<String> connections = failureAndSuccessConnection().collect(Collectors.toSet()); |
|||
defaultTbContext.tellNext(msg, connections); |
|||
|
|||
// THEN
|
|||
then(nodeCtxMock).should().getChainActor(); |
|||
then(nodeCtxMock).shouldHaveNoMoreInteractions(); |
|||
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(), nullable(Throwable.class), nullable(String.class)); |
|||
then(mainCtxMock).shouldHaveNoMoreInteractions(); |
|||
assertThat(nodeConnectionsCaptor.getAllValues()).hasSize(wantedNumberOfInvocations); |
|||
assertThat(nodeConnectionsCaptor.getAllValues()).containsExactlyInAnyOrderElementsOf(connections); |
|||
checkTellNextCommonLogic(callbackMock, connections, msg); |
|||
} |
|||
|
|||
private static Stream<String> failureAndSuccessConnection() { |
|||
return Stream.of(TbNodeConnectionType.FAILURE, TbNodeConnectionType.SUCCESS); |
|||
} |
|||
|
|||
@Test |
|||
public void givenDebugFailuresEventsAndFailureConnection_whenOutput_thenVerifyDebugOutputPersisted() { |
|||
// GIVEN
|
|||
var msgMock = mock(TbMsg.class); |
|||
var ruleNode = new RuleNode(RULE_NODE_ID); |
|||
ruleNode.setRuleChainId(RULE_CHAIN_ID); |
|||
ruleNode.setDebugSettings(DebugSettings.failures()); |
|||
given(msgMock.popFormStack()).willReturn(new TbMsgProcessingStackItem(RULE_CHAIN_ID, RULE_NODE_ID)); |
|||
given(nodeCtxMock.getTenantId()).willReturn(TENANT_ID); |
|||
given(nodeCtxMock.getSelf()).willReturn(ruleNode); |
|||
given(nodeCtxMock.getChainActor()).willReturn(chainActorMock); |
|||
|
|||
// WHEN
|
|||
defaultTbContext.output(msgMock, TbNodeConnectionType.FAILURE); |
|||
|
|||
// THEN
|
|||
checkOutputCommonLogic(msgMock, TbNodeConnectionType.FAILURE); |
|||
then(mainCtxMock).should().persistDebugOutput(TENANT_ID, RULE_NODE_ID, msgMock, TbNodeConnectionType.FAILURE, null, null); |
|||
then(mainCtxMock).shouldHaveNoMoreInteractions(); |
|||
then(nodeCtxMock).shouldHaveNoMoreInteractions(); |
|||
} |
|||
|
|||
@Test |
|||
public void givenDebugFailuresEventsAndSuccessConnection_whenOutput_thenVerifyDebugOutputNotPersisted() { |
|||
// GIVEN
|
|||
var msgMock = mock(TbMsg.class); |
|||
var ruleNode = new RuleNode(RULE_NODE_ID); |
|||
ruleNode.setRuleChainId(RULE_CHAIN_ID); |
|||
ruleNode.setDebugSettings(DebugSettings.failures()); |
|||
given(msgMock.popFormStack()).willReturn(new TbMsgProcessingStackItem(RULE_CHAIN_ID, RULE_NODE_ID)); |
|||
given(nodeCtxMock.getSelf()).willReturn(ruleNode); |
|||
given(nodeCtxMock.getChainActor()).willReturn(chainActorMock); |
|||
|
|||
// WHEN
|
|||
defaultTbContext.output(msgMock, TbNodeConnectionType.SUCCESS); |
|||
|
|||
// THEN
|
|||
checkOutputCommonLogic(msgMock, TbNodeConnectionType.SUCCESS); |
|||
then(mainCtxMock).shouldHaveNoMoreInteractions(); |
|||
then(nodeCtxMock).shouldHaveNoMoreInteractions(); |
|||
} |
|||
|
|||
@ParameterizedTest |
|||
@ValueSource(strings = {TbNodeConnectionType.SUCCESS, TbNodeConnectionType.FAILURE}) |
|||
void givenDebugDisabled_whenOutput_thenVerifyDebugOutputNotPersisted(String nodeConnection) { |
|||
// GIVEN
|
|||
var msgMock = mock(TbMsg.class); |
|||
var ruleNode = new RuleNode(RULE_NODE_ID); |
|||
ruleNode.setRuleChainId(RULE_CHAIN_ID); |
|||
given(msgMock.popFormStack()).willReturn(new TbMsgProcessingStackItem(RULE_CHAIN_ID, RULE_NODE_ID)); |
|||
given(nodeCtxMock.getSelf()).willReturn(ruleNode); |
|||
given(nodeCtxMock.getChainActor()).willReturn(chainActorMock); |
|||
|
|||
// WHEN
|
|||
defaultTbContext.output(msgMock, nodeConnection); |
|||
|
|||
// THEN
|
|||
checkOutputCommonLogic(msgMock, nodeConnection); |
|||
then(mainCtxMock).shouldHaveNoMoreInteractions(); |
|||
then(nodeCtxMock).shouldHaveNoMoreInteractions(); |
|||
} |
|||
|
|||
@ParameterizedTest |
|||
@ValueSource(strings = {TbNodeConnectionType.SUCCESS, TbNodeConnectionType.FAILURE}) |
|||
void givenDebugAllEvents_whenOutput_thenVerifyDebugOutputPersisted(String nodeConnection) { |
|||
// GIVEN
|
|||
var msgMock = mock(TbMsg.class); |
|||
var ruleNode = new RuleNode(RULE_NODE_ID); |
|||
ruleNode.setRuleChainId(RULE_CHAIN_ID); |
|||
ruleNode.setDebugSettings(DebugSettings.until(getUntilTime())); |
|||
given(msgMock.popFormStack()).willReturn(new TbMsgProcessingStackItem(RULE_CHAIN_ID, RULE_NODE_ID)); |
|||
given(nodeCtxMock.getTenantId()).willReturn(TENANT_ID); |
|||
given(nodeCtxMock.getSelf()).willReturn(ruleNode); |
|||
given(nodeCtxMock.getChainActor()).willReturn(chainActorMock); |
|||
|
|||
// WHEN
|
|||
defaultTbContext.output(msgMock, nodeConnection); |
|||
|
|||
// THEN
|
|||
checkOutputCommonLogic(msgMock, nodeConnection); |
|||
then(mainCtxMock).should().persistDebugOutput(TENANT_ID, RULE_NODE_ID, msgMock, nodeConnection, null, null); |
|||
then(mainCtxMock).shouldHaveNoMoreInteractions(); |
|||
then(nodeCtxMock).shouldHaveNoMoreInteractions(); |
|||
} |
|||
|
|||
@ParameterizedTest |
|||
@ValueSource(strings = {TbNodeConnectionType.SUCCESS, TbNodeConnectionType.FAILURE}) |
|||
void givenDebugAllThenOnlyFailureEvents_whenOutput_thenVerifyDebugOutputPersisted(String nodeConnection) { |
|||
// GIVEN
|
|||
var msgMock = mock(TbMsg.class); |
|||
var ruleNode = new RuleNode(RULE_NODE_ID); |
|||
ruleNode.setRuleChainId(RULE_CHAIN_ID); |
|||
ruleNode.setDebugSettings(DebugSettings.until(getUntilTime())); |
|||
given(msgMock.popFormStack()).willReturn(new TbMsgProcessingStackItem(RULE_CHAIN_ID, RULE_NODE_ID)); |
|||
given(nodeCtxMock.getTenantId()).willReturn(TENANT_ID); |
|||
given(nodeCtxMock.getSelf()).willReturn(ruleNode); |
|||
given(nodeCtxMock.getChainActor()).willReturn(chainActorMock); |
|||
|
|||
// WHEN
|
|||
defaultTbContext.output(msgMock, nodeConnection); |
|||
|
|||
// THEN
|
|||
checkOutputCommonLogic(msgMock, nodeConnection); |
|||
then(mainCtxMock).should().persistDebugOutput(TENANT_ID, RULE_NODE_ID, msgMock, nodeConnection, null, null); |
|||
then(mainCtxMock).shouldHaveNoMoreInteractions(); |
|||
then(nodeCtxMock).shouldHaveNoMoreInteractions(); |
|||
} |
|||
|
|||
@Test |
|||
public void givenEmptyStack_whenOutput_thenVerifyMsgAck() { |
|||
// GIVEN
|
|||
var msgMock = mock(TbMsg.class); |
|||
var ruleNode = new RuleNode(RULE_NODE_ID); |
|||
ruleNode.setRuleChainId(RULE_CHAIN_ID); |
|||
given(msgMock.popFormStack()).willReturn(null); |
|||
TbMsgCallback callbackMock = mock(TbMsgCallback.class); |
|||
given(msgMock.getCallback()).willReturn(callbackMock); |
|||
given(nodeCtxMock.getSelf()).willReturn(ruleNode); |
|||
|
|||
// WHEN
|
|||
defaultTbContext.output(msgMock, TbNodeConnectionType.SUCCESS); |
|||
|
|||
// THEN
|
|||
then(msgMock).should().popFormStack(); |
|||
then(callbackMock).should().onProcessingEnd(RULE_NODE_ID); |
|||
then(callbackMock).should().onSuccess(); |
|||
then(nodeCtxMock).should(never()).getChainActor(); |
|||
} |
|||
|
|||
@Test |
|||
public void givenEmptyStackAndDebugAllEvents_whenOutput_thenVerifyMsgAckAndDebugOutputPersisted() { |
|||
// GIVEN
|
|||
var msgMock = mock(TbMsg.class); |
|||
var ruleNode = new RuleNode(RULE_NODE_ID); |
|||
ruleNode.setRuleChainId(RULE_CHAIN_ID); |
|||
ruleNode.setDebugSettings(DebugSettings.until(getUntilTime())); |
|||
given(msgMock.popFormStack()).willReturn(null); |
|||
TbMsgCallback callbackMock = mock(TbMsgCallback.class); |
|||
given(msgMock.getCallback()).willReturn(callbackMock); |
|||
given(nodeCtxMock.getSelf()).willReturn(ruleNode); |
|||
given(nodeCtxMock.getTenantId()).willReturn(TENANT_ID); |
|||
|
|||
// WHEN
|
|||
defaultTbContext.output(msgMock, TbNodeConnectionType.SUCCESS); |
|||
|
|||
// THEN
|
|||
then(msgMock).should().popFormStack(); |
|||
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, null, null); |
|||
} |
|||
|
|||
@Test |
|||
public void givenEmptyStackAndDebugAllThenOnlyFailureEvents_whenOutput_thenVerifyMsgAckAndDebugOutputPersisted() { |
|||
// GIVEN
|
|||
var msgMock = mock(TbMsg.class); |
|||
var ruleNode = new RuleNode(RULE_NODE_ID); |
|||
ruleNode.setRuleChainId(RULE_CHAIN_ID); |
|||
ruleNode.setDebugSettings(DebugSettings.failuresOrUntil(getUntilTime())); |
|||
given(msgMock.popFormStack()).willReturn(null); |
|||
TbMsgCallback callbackMock = mock(TbMsgCallback.class); |
|||
given(msgMock.getCallback()).willReturn(callbackMock); |
|||
given(nodeCtxMock.getSelf()).willReturn(ruleNode); |
|||
given(nodeCtxMock.getTenantId()).willReturn(TENANT_ID); |
|||
|
|||
// WHEN
|
|||
defaultTbContext.output(msgMock, TbNodeConnectionType.SUCCESS); |
|||
|
|||
// THEN
|
|||
then(msgMock).should().popFormStack(); |
|||
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, null, null); |
|||
} |
|||
|
|||
@Test |
|||
public void givenDebugFailuresEvents_whenEnqueueForTellFailure_thenVerifyDebugOutputPersisted() { |
|||
// GIVEN
|
|||
var msg = getTbMsgWithQueueName(); |
|||
var tpi = new TopicPartitionInfo(DataConstants.MAIN_QUEUE_TOPIC, TENANT_ID, 0, true); |
|||
var ruleNode = new RuleNode(RULE_NODE_ID); |
|||
ruleNode.setRuleChainId(RULE_CHAIN_ID); |
|||
ruleNode.setDebugSettings(DebugSettings.failures()); |
|||
var tbClusterServiceMock = mock(TbClusterService.class); |
|||
|
|||
given(nodeCtxMock.getTenantId()).willReturn(TENANT_ID); |
|||
given(nodeCtxMock.getSelf()).willReturn(ruleNode); |
|||
given(mainCtxMock.resolve(any(ServiceType.class), anyString(), any(TenantId.class), any(EntityId.class))).willReturn(tpi); |
|||
given(mainCtxMock.getClusterService()).willReturn(tbClusterServiceMock); |
|||
|
|||
// WHEN
|
|||
defaultTbContext.enqueueForTellFailure(msg, EXCEPTION); |
|||
|
|||
// THEN
|
|||
then(mainCtxMock).should().resolve(ServiceType.TB_RULE_ENGINE, DataConstants.MAIN_QUEUE_NAME, TENANT_ID, TENANT_ID); |
|||
TbMsg expectedTbMsg = TbMsg.newMsg(msg, msg.getQueueName(), RULE_CHAIN_ID, RULE_NODE_ID); |
|||
checkEnqueueForTellFailurePushMsgToRuleEngine(tbClusterServiceMock, tpi, expectedTbMsg); |
|||
ArgumentCaptor<TbMsg> tbMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|||
then(mainCtxMock).should().persistDebugOutput(eq(TENANT_ID), eq(RULE_NODE_ID), tbMsgCaptor.capture(), eq(TbNodeConnectionType.FAILURE), isNull(), eq(EXCEPTION_MSG)); |
|||
TbMsg actualTbMsg = tbMsgCaptor.getValue(); |
|||
assertThat(actualTbMsg).usingRecursiveComparison() |
|||
.ignoringFields("id", "ctx") |
|||
.isEqualTo(expectedTbMsg); |
|||
then(mainCtxMock).should().getClusterService(); |
|||
then(mainCtxMock).shouldHaveNoMoreInteractions(); |
|||
then(tbClusterServiceMock).shouldHaveNoMoreInteractions(); |
|||
} |
|||
|
|||
@Test |
|||
public void givenDebugDisabled_whenEnqueueForTellFailure_thenVerifyDebugOutputNotPersisted() { |
|||
// GIVEN
|
|||
var msg = getTbMsgWithQueueName(); |
|||
var tpi = new TopicPartitionInfo(DataConstants.MAIN_QUEUE_TOPIC, TENANT_ID, 0, true); |
|||
var ruleNode = new RuleNode(RULE_NODE_ID); |
|||
ruleNode.setRuleChainId(RULE_CHAIN_ID); |
|||
var tbClusterServiceMock = mock(TbClusterService.class); |
|||
|
|||
given(nodeCtxMock.getTenantId()).willReturn(TENANT_ID); |
|||
given(nodeCtxMock.getSelf()).willReturn(ruleNode); |
|||
given(mainCtxMock.resolve(any(ServiceType.class), anyString(), any(TenantId.class), any(EntityId.class))).willReturn(tpi); |
|||
given(mainCtxMock.getClusterService()).willReturn(tbClusterServiceMock); |
|||
|
|||
// WHEN
|
|||
defaultTbContext.enqueueForTellFailure(msg, EXCEPTION); |
|||
|
|||
// THEN
|
|||
then(mainCtxMock).should().resolve(ServiceType.TB_RULE_ENGINE, DataConstants.MAIN_QUEUE_NAME, TENANT_ID, TENANT_ID); |
|||
TbMsg expectedTbMsg = TbMsg.newMsg(msg, msg.getQueueName(), RULE_CHAIN_ID, RULE_NODE_ID); |
|||
checkEnqueueForTellFailurePushMsgToRuleEngine(tbClusterServiceMock, tpi, expectedTbMsg); |
|||
then(mainCtxMock).should().getClusterService(); |
|||
then(mainCtxMock).shouldHaveNoMoreInteractions(); |
|||
then(tbClusterServiceMock).shouldHaveNoMoreInteractions(); |
|||
} |
|||
|
|||
@Test |
|||
public void givenDebugAllEvents_whenEnqueueForTellFailure_thenVerifyDebugOutputPersisted() { |
|||
// GIVEN
|
|||
var msg = getTbMsgWithQueueName(); |
|||
var tpi = new TopicPartitionInfo(DataConstants.MAIN_QUEUE_TOPIC, TENANT_ID, 0, true); |
|||
var ruleNode = new RuleNode(RULE_NODE_ID); |
|||
ruleNode.setRuleChainId(RULE_CHAIN_ID); |
|||
ruleNode.setDebugSettings(DebugSettings.until(getUntilTime())); |
|||
var tbClusterServiceMock = mock(TbClusterService.class); |
|||
|
|||
given(nodeCtxMock.getTenantId()).willReturn(TENANT_ID); |
|||
given(nodeCtxMock.getSelf()).willReturn(ruleNode); |
|||
given(mainCtxMock.resolve(any(ServiceType.class), anyString(), any(TenantId.class), any(EntityId.class))).willReturn(tpi); |
|||
given(mainCtxMock.getClusterService()).willReturn(tbClusterServiceMock); |
|||
|
|||
// WHEN
|
|||
defaultTbContext.enqueueForTellFailure(msg, EXCEPTION); |
|||
|
|||
// THEN
|
|||
then(mainCtxMock).should().resolve(ServiceType.TB_RULE_ENGINE, DataConstants.MAIN_QUEUE_NAME, TENANT_ID, TENANT_ID); |
|||
TbMsg expectedTbMsg = TbMsg.newMsg(msg, msg.getQueueName(), RULE_CHAIN_ID, RULE_NODE_ID); |
|||
checkEnqueueForTellFailurePushMsgToRuleEngine(tbClusterServiceMock, tpi, expectedTbMsg); |
|||
ArgumentCaptor<TbMsg> tbMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|||
then(mainCtxMock).should().persistDebugOutput(eq(TENANT_ID), eq(RULE_NODE_ID), tbMsgCaptor.capture(), eq(TbNodeConnectionType.FAILURE), isNull(), eq(EXCEPTION_MSG)); |
|||
TbMsg actualTbMsg = tbMsgCaptor.getValue(); |
|||
assertThat(actualTbMsg).usingRecursiveComparison() |
|||
.ignoringFields("id", "ctx") |
|||
.isEqualTo(expectedTbMsg); |
|||
then(mainCtxMock).should().getClusterService(); |
|||
then(mainCtxMock).shouldHaveNoMoreInteractions(); |
|||
then(tbClusterServiceMock).shouldHaveNoMoreInteractions(); |
|||
} |
|||
|
|||
@Test |
|||
public void givenInvalidMsg_whenEnqueueForTellFailure_thenDoNothing() { |
|||
// GIVEN
|
|||
var msgMock = mock(TbMsg.class); |
|||
var tpi = new TopicPartitionInfo(DataConstants.MAIN_QUEUE_TOPIC, TENANT_ID, 0, true); |
|||
|
|||
given(msgMock.getOriginator()).willReturn(TENANT_ID); |
|||
given(msgMock.getQueueName()).willReturn(DataConstants.MAIN_QUEUE_NAME); |
|||
given(msgMock.isValid()).willReturn(false); |
|||
given(nodeCtxMock.getTenantId()).willReturn(TENANT_ID); |
|||
given(mainCtxMock.resolve(any(ServiceType.class), anyString(), any(TenantId.class), any(EntityId.class))).willReturn(tpi); |
|||
|
|||
// WHEN
|
|||
defaultTbContext.enqueueForTellFailure(msgMock, EXCEPTION); |
|||
|
|||
// THEN
|
|||
then(msgMock).should(times(2)).getQueueName(); |
|||
then(msgMock).should().getOriginator(); |
|||
then(msgMock).should().isValid(); |
|||
then(msgMock).shouldHaveNoMoreInteractions(); |
|||
|
|||
then(mainCtxMock).should().resolve(ServiceType.TB_RULE_ENGINE, DataConstants.MAIN_QUEUE_NAME, TENANT_ID, TENANT_ID); |
|||
then(mainCtxMock).shouldHaveNoMoreInteractions(); |
|||
|
|||
then(nodeCtxMock).should(times(2)).getTenantId(); |
|||
then(nodeCtxMock).shouldHaveNoMoreInteractions(); |
|||
then(chainActorMock).shouldHaveNoInteractions(); |
|||
} |
|||
|
|||
@MethodSource |
|||
@ParameterizedTest |
|||
void givenDebugOptions_whenEnqueueForTellNext_thenVerifyDebugOutputPersistedOnlyForDebugAll(boolean debugFailures, long debugAllUntil, String connectionType) { |
|||
// GIVEN
|
|||
var msg = getTbMsgWithQueueName(); |
|||
var tpi = new TopicPartitionInfo(DataConstants.MAIN_QUEUE_TOPIC, TENANT_ID, 0, true); |
|||
var ruleNode = new RuleNode(RULE_NODE_ID); |
|||
ruleNode.setRuleChainId(RULE_CHAIN_ID); |
|||
ruleNode.setDebugSettings(new DebugSettings(debugFailures, debugAllUntil)); |
|||
var tbClusterServiceMock = mock(TbClusterService.class); |
|||
|
|||
given(nodeCtxMock.getTenantId()).willReturn(TENANT_ID); |
|||
given(nodeCtxMock.getSelf()).willReturn(ruleNode); |
|||
given(mainCtxMock.resolve(any(ServiceType.class), anyString(), any(TenantId.class), any(EntityId.class))).willReturn(tpi); |
|||
given(mainCtxMock.getClusterService()).willReturn(tbClusterServiceMock); |
|||
|
|||
// WHEN
|
|||
defaultTbContext.enqueueForTellNext(msg, connectionType); |
|||
|
|||
// THEN
|
|||
then(mainCtxMock).should().resolve(ServiceType.TB_RULE_ENGINE, DataConstants.MAIN_QUEUE_NAME, TENANT_ID, TENANT_ID); |
|||
TbMsg expectedTbMsg = TbMsg.newMsg(msg, msg.getQueueName(), RULE_CHAIN_ID, RULE_NODE_ID); |
|||
|
|||
ArgumentCaptor<ToRuleEngineMsg> toRuleEngineMsgCaptor = ArgumentCaptor.forClass(ToRuleEngineMsg.class); |
|||
ArgumentCaptor<SimpleTbQueueCallback> simpleTbQueueCallbackCaptor = ArgumentCaptor.forClass(SimpleTbQueueCallback.class); |
|||
then(tbClusterServiceMock).should().pushMsgToRuleEngine(eq(tpi), notNull(UUID.class), toRuleEngineMsgCaptor.capture(), simpleTbQueueCallbackCaptor.capture()); |
|||
|
|||
ToRuleEngineMsg actualToRuleEngineMsg = toRuleEngineMsgCaptor.getValue(); |
|||
assertThat(actualToRuleEngineMsg).usingRecursiveComparison() |
|||
.ignoringFields("tbMsg_") |
|||
.isEqualTo(ToRuleEngineMsg.newBuilder() |
|||
.setTenantIdMSB(TENANT_ID.getId().getMostSignificantBits()) |
|||
.setTenantIdLSB(TENANT_ID.getId().getLeastSignificantBits()) |
|||
.setTbMsg(TbMsg.toByteString(expectedTbMsg)) |
|||
.addAllRelationTypes(List.of(connectionType)).build()); |
|||
|
|||
var simpleTbQueueCallback = simpleTbQueueCallbackCaptor.getValue(); |
|||
assertThat(simpleTbQueueCallback).isNotNull(); |
|||
simpleTbQueueCallback.onSuccess(null); |
|||
|
|||
if (debugAllUntil > 0) { |
|||
ArgumentCaptor<TbMsg> tbMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|||
then(mainCtxMock).should().persistDebugOutput(eq(TENANT_ID), eq(RULE_NODE_ID), tbMsgCaptor.capture(), eq(connectionType), isNull(), isNull()); |
|||
TbMsg actualTbMsg = tbMsgCaptor.getValue(); |
|||
assertThat(actualTbMsg).usingRecursiveComparison() |
|||
.ignoringFields("id", "ctx") |
|||
.isEqualTo(expectedTbMsg); |
|||
} |
|||
then(mainCtxMock).should().getClusterService(); |
|||
then(mainCtxMock).shouldHaveNoMoreInteractions(); |
|||
then(tbClusterServiceMock).shouldHaveNoMoreInteractions(); |
|||
} |
|||
|
|||
@MethodSource |
|||
@ParameterizedTest |
|||
void givenDebugOptions_whenEnqueue_thenVerifyDebugOutputPersistedOnlyForDebugAll(boolean debugFailures, long debugAllUntil) { |
|||
// GIVEN
|
|||
var msg = getTbMsgWithQueueName(); |
|||
var tpi = new TopicPartitionInfo(DataConstants.MAIN_QUEUE_TOPIC, TENANT_ID, 0, true); |
|||
var ruleNode = new RuleNode(RULE_NODE_ID); |
|||
ruleNode.setQueueName(DataConstants.MAIN_QUEUE_NAME); |
|||
ruleNode.setRuleChainId(RULE_CHAIN_ID); |
|||
ruleNode.setDebugSettings(new DebugSettings(debugFailures, debugAllUntil)); |
|||
var tbClusterServiceMock = mock(TbClusterService.class); |
|||
|
|||
given(nodeCtxMock.getTenantId()).willReturn(TENANT_ID); |
|||
given(nodeCtxMock.getSelf()).willReturn(ruleNode); |
|||
given(mainCtxMock.resolve(any(ServiceType.class), anyString(), any(TenantId.class), any(EntityId.class))).willReturn(tpi); |
|||
given(mainCtxMock.getClusterService()).willReturn(tbClusterServiceMock); |
|||
|
|||
Consumer<Throwable> onFailure = mock(Consumer.class); |
|||
Runnable onSuccess = mock(Runnable.class); |
|||
|
|||
// WHEN
|
|||
defaultTbContext.enqueue(msg, onSuccess, onFailure); |
|||
|
|||
// THEN
|
|||
then(mainCtxMock).should().resolve(ServiceType.TB_RULE_ENGINE, DataConstants.MAIN_QUEUE_NAME, TENANT_ID, TENANT_ID); |
|||
TbMsg expectedTbMsg = TbMsg.newMsg(msg, msg.getQueueName(), RULE_CHAIN_ID, RULE_NODE_ID); |
|||
|
|||
ArgumentCaptor<ToRuleEngineMsg> toRuleEngineMsgCaptor = ArgumentCaptor.forClass(ToRuleEngineMsg.class); |
|||
ArgumentCaptor<SimpleTbQueueCallback> simpleTbQueueCallbackCaptor = ArgumentCaptor.forClass(SimpleTbQueueCallback.class); |
|||
then(tbClusterServiceMock).should().pushMsgToRuleEngine(eq(tpi), notNull(UUID.class), toRuleEngineMsgCaptor.capture(), simpleTbQueueCallbackCaptor.capture()); |
|||
|
|||
ToRuleEngineMsg actualToRuleEngineMsg = toRuleEngineMsgCaptor.getValue(); |
|||
assertThat(actualToRuleEngineMsg).usingRecursiveComparison() |
|||
.ignoringFields("tbMsg_") |
|||
.isEqualTo(ToRuleEngineMsg.newBuilder() |
|||
.setTenantIdMSB(TENANT_ID.getId().getMostSignificantBits()) |
|||
.setTenantIdLSB(TENANT_ID.getId().getLeastSignificantBits()) |
|||
.setTbMsg(TbMsg.toByteString(expectedTbMsg)) |
|||
.build()); |
|||
|
|||
var simpleTbQueueCallback = simpleTbQueueCallbackCaptor.getValue(); |
|||
assertThat(simpleTbQueueCallback).isNotNull(); |
|||
simpleTbQueueCallback.onSuccess(null); |
|||
|
|||
if (debugAllUntil > 0) { |
|||
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(); |
|||
then(tbClusterServiceMock).shouldHaveNoMoreInteractions(); |
|||
} |
|||
|
|||
@Test |
|||
public void givenDebugFailuress_whenTellFailure_thenVerifyDebugOutputPersisted() { |
|||
// GIVEN
|
|||
var msg = getTbMsg(); |
|||
var ruleNode = new RuleNode(RULE_NODE_ID); |
|||
ruleNode.setRuleChainId(RULE_CHAIN_ID); |
|||
ruleNode.setDebugSettings(DebugSettings.failures()); |
|||
given(nodeCtxMock.getTenantId()).willReturn(TENANT_ID); |
|||
given(nodeCtxMock.getSelf()).willReturn(ruleNode); |
|||
given(nodeCtxMock.getChainActor()).willReturn(chainActorMock); |
|||
|
|||
// WHEN
|
|||
defaultTbContext.tellFailure(msg, EXCEPTION); |
|||
|
|||
// THEN
|
|||
var expectedRuleNodeToRuleChainTellNextMsg = new RuleNodeToRuleChainTellNextMsg( |
|||
RULE_CHAIN_ID, |
|||
RULE_NODE_ID, |
|||
Collections.singleton(TbNodeConnectionType.FAILURE), |
|||
msg, |
|||
EXCEPTION_MSG |
|||
); |
|||
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, null); |
|||
then(mainCtxMock).shouldHaveNoMoreInteractions(); |
|||
then(nodeCtxMock).shouldHaveNoMoreInteractions(); |
|||
} |
|||
|
|||
@Test |
|||
public void givenDebugDisabled_whenTellFailure_thenVerifyDebugOutputNotPersisted() { |
|||
// GIVEN
|
|||
var msg = getTbMsg(); |
|||
var ruleNode = new RuleNode(RULE_NODE_ID); |
|||
ruleNode.setRuleChainId(RULE_CHAIN_ID); |
|||
given(nodeCtxMock.getSelf()).willReturn(ruleNode); |
|||
given(nodeCtxMock.getChainActor()).willReturn(chainActorMock); |
|||
|
|||
// WHEN
|
|||
defaultTbContext.tellFailure(msg, EXCEPTION); |
|||
|
|||
// THEN
|
|||
var expectedRuleNodeToRuleChainTellNextMsg = new RuleNodeToRuleChainTellNextMsg( |
|||
RULE_CHAIN_ID, |
|||
RULE_NODE_ID, |
|||
Collections.singleton(TbNodeConnectionType.FAILURE), |
|||
msg, |
|||
EXCEPTION_MSG |
|||
); |
|||
then(chainActorMock).should().tell(expectedRuleNodeToRuleChainTellNextMsg); |
|||
then(chainActorMock).shouldHaveNoMoreInteractions(); |
|||
then(nodeCtxMock).should().getChainActor(); |
|||
then(mainCtxMock).shouldHaveNoInteractions(); |
|||
then(nodeCtxMock).shouldHaveNoMoreInteractions(); |
|||
} |
|||
|
|||
@Test |
|||
public void givenDebugAllEvents_whenTellFailure_thenVerifyDebugOutputPersisted() { |
|||
// GIVEN
|
|||
var msg = getTbMsg(); |
|||
var ruleNode = new RuleNode(RULE_NODE_ID); |
|||
ruleNode.setRuleChainId(RULE_CHAIN_ID); |
|||
ruleNode.setDebugSettings(DebugSettings.until(getUntilTime())); |
|||
given(nodeCtxMock.getTenantId()).willReturn(TENANT_ID); |
|||
given(nodeCtxMock.getSelf()).willReturn(ruleNode); |
|||
given(nodeCtxMock.getChainActor()).willReturn(chainActorMock); |
|||
|
|||
// WHEN
|
|||
defaultTbContext.tellFailure(msg, EXCEPTION); |
|||
|
|||
// THEN
|
|||
var expectedRuleNodeToRuleChainTellNextMsg = new RuleNodeToRuleChainTellNextMsg( |
|||
RULE_CHAIN_ID, |
|||
RULE_NODE_ID, |
|||
Collections.singleton(TbNodeConnectionType.FAILURE), |
|||
msg, |
|||
EXCEPTION_MSG |
|||
); |
|||
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, null); |
|||
then(mainCtxMock).shouldHaveNoMoreInteractions(); |
|||
then(nodeCtxMock).shouldHaveNoMoreInteractions(); |
|||
} |
|||
|
|||
@MethodSource |
|||
@ParameterizedTest |
|||
void givenDebugFailuresAndDebugAllAndConnectionAndPersistedResultOptions_whenTellNext_thenVerifyDebugOutputPersistence(boolean debugFailures, |
|||
long debugAllUntil, |
|||
String connection, |
|||
boolean shouldPersist, |
|||
boolean shouldPersistAfterDurationTime) { |
|||
// GIVEN
|
|||
var callbackMock = mock(TbMsgCallback.class); |
|||
var msg = getTbMsgWithCallback(callbackMock); |
|||
var ruleNode = new RuleNode(RULE_NODE_ID); |
|||
ruleNode.setRuleChainId(RULE_CHAIN_ID); |
|||
ruleNode.setDebugSettings(new DebugSettings(debugFailures, debugAllUntil)); |
|||
if (shouldPersist) { |
|||
given(nodeCtxMock.getTenantId()).willReturn(TENANT_ID); |
|||
} |
|||
given(nodeCtxMock.getSelf()).willReturn(ruleNode); |
|||
given(nodeCtxMock.getChainActor()).willReturn(chainActorMock); |
|||
|
|||
// WHEN
|
|||
defaultTbContext.tellNext(msg, connection); |
|||
|
|||
// THEN
|
|||
if (shouldPersist) { |
|||
then(mainCtxMock).should().persistDebugOutput(TENANT_ID, RULE_NODE_ID, msg, connection, null, null); |
|||
} |
|||
|
|||
// GIVEN
|
|||
Mockito.clearInvocations(mainCtxMock); |
|||
ruleNode.setDebugSettings(new DebugSettings(ruleNode.getDebugSettings().isFailuresEnabled(), 0)); |
|||
|
|||
// WHEN
|
|||
defaultTbContext.tellNext(msg, connection); |
|||
|
|||
// THEN
|
|||
if (shouldPersistAfterDurationTime) { |
|||
then(mainCtxMock).should().persistDebugOutput(TENANT_ID, RULE_NODE_ID, msg, connection, null, null); |
|||
} |
|||
} |
|||
|
|||
private void checkTellNextCommonLogic(TbMsgCallback callbackMock, String nodeConnection, TbMsg msg) { |
|||
checkTellNextCommonLogic(callbackMock, Collections.singleton(nodeConnection), msg); |
|||
} |
|||
|
|||
private void checkTellNextCommonLogic(TbMsgCallback callbackMock, Set<String> nodeConnections, TbMsg msg) { |
|||
then(callbackMock).should().onProcessingEnd(RULE_NODE_ID); |
|||
then(callbackMock).shouldHaveNoMoreInteractions(); |
|||
var expectedRuleNodeToRuleChainTellNextMsg = new RuleNodeToRuleChainTellNextMsg( |
|||
RULE_CHAIN_ID, |
|||
RULE_NODE_ID, |
|||
nodeConnections, |
|||
msg, |
|||
null); |
|||
then(chainActorMock).should().tell(expectedRuleNodeToRuleChainTellNextMsg); |
|||
then(chainActorMock).shouldHaveNoMoreInteractions(); |
|||
} |
|||
|
|||
private void checkOutputCommonLogic(TbMsg msg, String nodeConnection) { |
|||
then(msg).should().popFormStack(); |
|||
var expectedRuleChainOutputMsg = new RuleChainOutputMsg( |
|||
RULE_CHAIN_ID, |
|||
RULE_NODE_ID, |
|||
nodeConnection, |
|||
msg); |
|||
then(chainActorMock).should().tell(expectedRuleChainOutputMsg); |
|||
then(chainActorMock).shouldHaveNoMoreInteractions(); |
|||
then(nodeCtxMock).should().getChainActor(); |
|||
} |
|||
|
|||
private void checkEnqueueForTellFailurePushMsgToRuleEngine(TbClusterService tbClusterService, TopicPartitionInfo tpi, TbMsg expectedTbMsg) { |
|||
ArgumentCaptor<ToRuleEngineMsg> toRuleEngineMsgCaptor = ArgumentCaptor.forClass(ToRuleEngineMsg.class); |
|||
ArgumentCaptor<SimpleTbQueueCallback> simpleTbQueueCallbackCaptor = ArgumentCaptor.forClass(SimpleTbQueueCallback.class); |
|||
then(tbClusterService).should().pushMsgToRuleEngine(eq(tpi), notNull(UUID.class), toRuleEngineMsgCaptor.capture(), simpleTbQueueCallbackCaptor.capture()); |
|||
|
|||
ToRuleEngineMsg actualToRuleEngineMsg = toRuleEngineMsgCaptor.getValue(); |
|||
assertThat(actualToRuleEngineMsg).usingRecursiveComparison() |
|||
.ignoringFields("tbMsg_") |
|||
.isEqualTo(ToRuleEngineMsg.newBuilder() |
|||
.setTenantIdMSB(TENANT_ID.getId().getMostSignificantBits()) |
|||
.setTenantIdLSB(TENANT_ID.getId().getLeastSignificantBits()) |
|||
.setTbMsg(TbMsg.toByteString(expectedTbMsg)) |
|||
.setFailureMessage(EXCEPTION_MSG) |
|||
.addAllRelationTypes(List.of(TbNodeConnectionType.FAILURE)).build()); |
|||
|
|||
var simpleTbQueueCallback = simpleTbQueueCallbackCaptor.getValue(); |
|||
assertThat(simpleTbQueueCallback).isNotNull(); |
|||
simpleTbQueueCallback.onSuccess(null); |
|||
} |
|||
|
|||
private static Stream<Arguments> givenDebugOptions_whenEnqueueForTellNext_thenVerifyDebugOutputPersistedOnlyForDebugAll() { |
|||
return Stream.of( |
|||
Arguments.of(false, getUntilTime(), TbNodeConnectionType.OTHER), |
|||
Arguments.of(true, getUntilTime(), TbNodeConnectionType.OTHER), |
|||
Arguments.of(true, 0, TbNodeConnectionType.TRUE), |
|||
Arguments.of(false, 0, TbNodeConnectionType.FALSE) |
|||
); |
|||
} |
|||
|
|||
private static Stream<Arguments> givenDebugOptions_whenEnqueue_thenVerifyDebugOutputPersistedOnlyForDebugAll() { |
|||
return Stream.of( |
|||
Arguments.of(false, getUntilTime()), |
|||
Arguments.of(true, getUntilTime()), |
|||
Arguments.of(true, 0), |
|||
Arguments.of(false, 0) |
|||
); |
|||
} |
|||
|
|||
private static Stream<Arguments> givenDebugFailuresAndDebugAllAndConnectionAndPersistedResultOptions_whenTellNext_thenVerifyDebugOutputPersistence() { |
|||
return Stream.of( |
|||
Arguments.of(false, getUntilTime(), TbNodeConnectionType.SUCCESS, true, false), |
|||
Arguments.of(false, getUntilTime(), TbNodeConnectionType.FAILURE, true, false), |
|||
Arguments.of(true, getUntilTime(), TbNodeConnectionType.SUCCESS, true, false), |
|||
Arguments.of(true, getUntilTime(), TbNodeConnectionType.FAILURE, true, true), |
|||
Arguments.of(true, 0, TbNodeConnectionType.SUCCESS, false, false), |
|||
Arguments.of(true, 0, TbNodeConnectionType.FAILURE, true, true), |
|||
Arguments.of(false, 0, TbNodeConnectionType.SUCCESS, false, false), |
|||
Arguments.of(false, 0, TbNodeConnectionType.FAILURE, false, false) |
|||
); |
|||
} |
|||
|
|||
private TbMsg getTbMsgWithCallback(TbMsgCallback callback) { |
|||
return TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, TENANT_ID, TbMsgMetaData.EMPTY, TbMsg.EMPTY_STRING, callback); |
|||
} |
|||
|
|||
private TbMsg getTbMsgWithQueueName() { |
|||
return TbMsg.newMsg(DataConstants.MAIN_QUEUE_NAME, TbMsgType.POST_TELEMETRY_REQUEST, TENANT_ID, TbMsgMetaData.EMPTY, TbMsg.EMPTY_STRING); |
|||
} |
|||
|
|||
private TbMsg getTbMsg() { |
|||
return TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, TENANT_ID, TbMsgMetaData.EMPTY, TbMsg.EMPTY_STRING); |
|||
} |
|||
|
|||
private static long getUntilTime() { |
|||
return getUntilTime(15); |
|||
} |
|||
|
|||
private static long getUntilTime(int maxRuleNodeDebugModeDurationMinutes) { |
|||
return System.currentTimeMillis() + TimeUnit.MINUTES.toMillis(maxRuleNodeDebugModeDurationMinutes); |
|||
} |
|||
} |
|||
@ -0,0 +1,32 @@ |
|||
/** |
|||
* Copyright © 2016-2024 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.common.data; |
|||
|
|||
import org.thingsboard.server.common.data.debug.DebugSettings; |
|||
|
|||
public interface HasDebugSettings { |
|||
|
|||
@Deprecated |
|||
boolean isDebugMode(); |
|||
|
|||
@Deprecated |
|||
void setDebugMode(boolean debugMode); |
|||
|
|||
DebugSettings getDebugSettings(); |
|||
|
|||
void setDebugSettings(DebugSettings debugSettings); |
|||
|
|||
} |
|||
@ -0,0 +1,59 @@ |
|||
/** |
|||
* Copyright © 2016-2024 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.common.data.debug; |
|||
|
|||
import io.swagger.v3.oas.annotations.media.Schema; |
|||
import lombok.Data; |
|||
import lombok.NoArgsConstructor; |
|||
|
|||
@Data |
|||
@NoArgsConstructor |
|||
public class DebugSettings { |
|||
|
|||
private static DebugSettings DEBUG_OFF = new DebugSettings(false, 0); |
|||
private static DebugSettings DEBUG_FAILURES = new DebugSettings(true, 0); |
|||
|
|||
public DebugSettings(boolean failuresEnabled, long allEnabledUntil) { |
|||
this.failuresEnabled = failuresEnabled; |
|||
this.allEnabled = false; |
|||
this.allEnabledUntil = allEnabledUntil; |
|||
} |
|||
|
|||
@Schema(description = "Debug failures. ", example = "false") |
|||
private boolean failuresEnabled; |
|||
@Schema(description = "Debug All. Used as a trigger for updating debugAllUntil.", example = "false") |
|||
private boolean allEnabled; |
|||
@Schema(description = "Timestamp of the end time for the processing debug events.") |
|||
private long allEnabledUntil; |
|||
|
|||
public static DebugSettings off() {return DebugSettings.DEBUG_OFF;} |
|||
|
|||
public static DebugSettings failures() {return DebugSettings.DEBUG_FAILURES;} |
|||
|
|||
public static DebugSettings until(long ts) {return new DebugSettings(false, ts);} |
|||
|
|||
public static DebugSettings failuresOrUntil(long ts) {return new DebugSettings(true, ts);} |
|||
|
|||
public static DebugSettings all() { |
|||
var ds = new DebugSettings(); |
|||
ds.setAllEnabled(true); |
|||
return ds; |
|||
} |
|||
|
|||
public DebugSettings copy(long maxDebugAllUntil) { |
|||
return new DebugSettings(failuresEnabled, allEnabled ? maxDebugAllUntil : Math.min(allEnabledUntil, maxDebugAllUntil)); |
|||
} |
|||
} |
|||
@ -0,0 +1,60 @@ |
|||
/** |
|||
* Copyright © 2016-2024 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.common.util; |
|||
|
|||
import org.thingsboard.server.common.data.HasDebugSettings; |
|||
import org.thingsboard.server.common.data.msg.TbNodeConnectionType; |
|||
|
|||
import java.util.Set; |
|||
|
|||
public final class DebugModeUtil { |
|||
|
|||
private static final int DEBUG_MODE_DEFAULT_DURATION_MINUTES = 15; |
|||
|
|||
private DebugModeUtil() { |
|||
} |
|||
|
|||
public static int getMaxDebugAllDuration(int tenantProfileDuration, int systemDefaultDuration) { |
|||
if (tenantProfileDuration > 0) { |
|||
return tenantProfileDuration; |
|||
} else { |
|||
return systemDefaultDuration > 0 ? systemDefaultDuration : DEBUG_MODE_DEFAULT_DURATION_MINUTES; |
|||
} |
|||
} |
|||
|
|||
public static boolean isDebugAllAvailable(HasDebugSettings debugSettingsAware) { |
|||
var debugSettings = debugSettingsAware.getDebugSettings(); |
|||
return debugSettings != null && debugSettings.getAllEnabledUntil() > System.currentTimeMillis(); |
|||
} |
|||
|
|||
public static boolean isDebugAvailable(HasDebugSettings debugSettingsAware, String nodeConnection) { |
|||
if (isDebugAllAvailable(debugSettingsAware)) { |
|||
return true; |
|||
} else { |
|||
var debugSettings = debugSettingsAware.getDebugSettings(); |
|||
return debugSettings != null && debugSettings.isFailuresEnabled() && TbNodeConnectionType.FAILURE.equals(nodeConnection); |
|||
} |
|||
} |
|||
|
|||
public static boolean isDebugFailuresAvailable(HasDebugSettings debugSettingsAware, Set<String> nodeConnections) { |
|||
if (isDebugAllAvailable(debugSettingsAware)) { |
|||
return true; |
|||
} else { |
|||
var debugSettings = debugSettingsAware.getDebugSettings(); |
|||
return debugSettings != null && nodeConnections != null && debugSettings.isFailuresEnabled() && nodeConnections.contains(TbNodeConnectionType.FAILURE); |
|||
} |
|||
} |
|||
} |
|||
@ -0,0 +1,32 @@ |
|||
<!-- |
|||
|
|||
Copyright © 2016-2024 The Thingsboard Authors |
|||
|
|||
Licensed under the Apache License, Version 2.0 (the "License"); |
|||
you may not use this file except in compliance with the License. |
|||
You may obtain a copy of the License at |
|||
|
|||
http://www.apache.org/licenses/LICENSE-2.0 |
|||
|
|||
Unless required by applicable law or agreed to in writing, software |
|||
distributed under the License is distributed on an "AS IS" BASIS, |
|||
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
See the License for the specific language governing permissions and |
|||
limitations under the License. |
|||
|
|||
--> |
|||
<button mat-stroked-button |
|||
class="tb-rounded-btn w-36 flex-1" |
|||
color="primary" |
|||
#matButton |
|||
[class.active]="((isDebugAllActive$ | async) || failuresEnabled) && !disabled" |
|||
[disabled]="disabled" |
|||
(click)="openDebugStrategyPanel($event, matButton)"> |
|||
<mat-icon [color]="debugSettingsFormGroup.disabled ? 'inherit' : 'primary'">bug_report</mat-icon> |
|||
<span *ngIf="!failuresEnabled && !(isDebugAllActive$ | async)" translate>common.disabled</span> |
|||
<span *ngIf="(isDebugAllActive$ | async) && failuresEnabled" translate>debug-config.all</span> |
|||
<span *ngIf="(isDebugAllActive$ | async) && !failuresEnabled"> |
|||
{{ !allEnabled ? (allEnabledUntil | durationLeft) : ('debug-config.min' | translate: { number: maxDebugModeDurationMinutes }) }} |
|||
</span> |
|||
<span *ngIf="!(isDebugAllActive$ | async) && failuresEnabled" translate>debug-config.failures</span> |
|||
</button> |
|||
@ -0,0 +1,155 @@ |
|||
///
|
|||
/// Copyright © 2016-2024 The Thingsboard Authors
|
|||
///
|
|||
/// Licensed under the Apache License, Version 2.0 (the "License");
|
|||
/// you may not use this file except in compliance with the License.
|
|||
/// You may obtain a copy of the License at
|
|||
///
|
|||
/// http://www.apache.org/licenses/LICENSE-2.0
|
|||
///
|
|||
/// Unless required by applicable law or agreed to in writing, software
|
|||
/// distributed under the License is distributed on an "AS IS" BASIS,
|
|||
/// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
|||
/// See the License for the specific language governing permissions and
|
|||
/// limitations under the License.
|
|||
///
|
|||
|
|||
import { ChangeDetectionStrategy, Component, forwardRef, Input, Renderer2, ViewContainerRef } from '@angular/core'; |
|||
import { CommonModule } from '@angular/common'; |
|||
import { SharedModule } from '@shared/shared.module'; |
|||
import { DurationLeftPipe } from '@shared/pipe/duration-left.pipe'; |
|||
import { TbPopoverService } from '@shared/components/popover.service'; |
|||
import { MatButton } from '@angular/material/button'; |
|||
import { DebugSettingsPanelComponent } from './debug-settings-panel.component'; |
|||
import { takeUntilDestroyed } from '@angular/core/rxjs-interop'; |
|||
import { of, shareReplay, timer } from 'rxjs'; |
|||
import { SECOND } from '@shared/models/time/time.models'; |
|||
import { DebugSettings } from '@shared/models/entity.models'; |
|||
import { map, startWith, switchMap, takeWhile } from 'rxjs/operators'; |
|||
import { getCurrentAuthState } from '@core/auth/auth.selectors'; |
|||
import { AppState } from '@core/core.state'; |
|||
import { Store } from '@ngrx/store'; |
|||
import { ControlValueAccessor, FormBuilder, NG_VALUE_ACCESSOR } from '@angular/forms'; |
|||
|
|||
@Component({ |
|||
selector: 'tb-debug-settings-button', |
|||
templateUrl: './debug-settings-button.component.html', |
|||
standalone: true, |
|||
imports: [ |
|||
CommonModule, |
|||
SharedModule, |
|||
DurationLeftPipe, |
|||
], |
|||
providers: [ |
|||
{ |
|||
provide: NG_VALUE_ACCESSOR, |
|||
useExisting: forwardRef(() => DebugSettingsButtonComponent), |
|||
multi: true |
|||
}, |
|||
], |
|||
changeDetection: ChangeDetectionStrategy.OnPush |
|||
}) |
|||
export class DebugSettingsButtonComponent implements ControlValueAccessor { |
|||
|
|||
@Input() debugLimitsConfiguration: string; |
|||
|
|||
debugSettingsFormGroup = this.fb.group({ |
|||
failuresEnabled: [false], |
|||
allEnabled: [false], |
|||
allEnabledUntil: [] |
|||
}); |
|||
|
|||
disabled = false; |
|||
|
|||
isDebugAllActive$ = this.debugSettingsFormGroup.get('allEnabled').valueChanges.pipe( |
|||
startWith(null), |
|||
switchMap(() => { |
|||
if (this.allEnabled) { |
|||
return of(true); |
|||
} else { |
|||
return timer(0, SECOND).pipe( |
|||
map(() => this.allEnabledUntil > new Date().getTime()), |
|||
takeWhile(value => value, true) |
|||
); |
|||
} |
|||
}), |
|||
takeUntilDestroyed(), |
|||
shareReplay(1) |
|||
); |
|||
|
|||
readonly maxDebugModeDurationMinutes = getCurrentAuthState(this.store).maxDebugModeDurationMinutes; |
|||
|
|||
private propagateChange: (settings: DebugSettings) => void = () => {}; |
|||
|
|||
constructor(private popoverService: TbPopoverService, |
|||
private renderer: Renderer2, |
|||
private store: Store<AppState>, |
|||
private viewContainerRef: ViewContainerRef, |
|||
private fb: FormBuilder, |
|||
) { |
|||
this.debugSettingsFormGroup.valueChanges.pipe( |
|||
takeUntilDestroyed() |
|||
).subscribe(value => { |
|||
this.propagateChange(value); |
|||
}) |
|||
} |
|||
|
|||
get failuresEnabled(): boolean { |
|||
return this.debugSettingsFormGroup.get('failuresEnabled').value; |
|||
} |
|||
|
|||
get allEnabled(): boolean { |
|||
return this.debugSettingsFormGroup.get('allEnabled').value; |
|||
} |
|||
|
|||
get allEnabledUntil(): number { |
|||
return this.debugSettingsFormGroup.get('allEnabledUntil').value; |
|||
} |
|||
|
|||
openDebugStrategyPanel($event: Event, matButton: MatButton): void { |
|||
if ($event) { |
|||
$event.stopPropagation(); |
|||
} |
|||
const trigger = matButton._elementRef.nativeElement; |
|||
const debugSettings = this.debugSettingsFormGroup.value; |
|||
|
|||
if (this.popoverService.hasPopover(trigger)) { |
|||
this.popoverService.hidePopover(trigger); |
|||
} else { |
|||
const debugStrategyPopover = this.popoverService.displayPopover(trigger, this.renderer, |
|||
this.viewContainerRef, DebugSettingsPanelComponent, 'bottom', true, null, |
|||
{ |
|||
...debugSettings, |
|||
maxDebugModeDurationMinutes: this.maxDebugModeDurationMinutes, |
|||
debugLimitsConfiguration: this.debugLimitsConfiguration |
|||
}, |
|||
{}, |
|||
{}, {}, true); |
|||
debugStrategyPopover.tbComponentRef.instance.popover = debugStrategyPopover; |
|||
debugStrategyPopover.tbComponentRef.instance.onSettingsApplied.subscribe((settings: DebugSettings) => { |
|||
this.debugSettingsFormGroup.patchValue(settings); |
|||
debugStrategyPopover.hide(); |
|||
}); |
|||
} |
|||
} |
|||
|
|||
registerOnChange(fn: (settings: DebugSettings) => void): void { |
|||
this.propagateChange = fn; |
|||
} |
|||
|
|||
registerOnTouched(_: () => void): void {} |
|||
|
|||
writeValue(settings: DebugSettings): void { |
|||
this.debugSettingsFormGroup.patchValue(settings, {emitEvent: false}); |
|||
this.debugSettingsFormGroup.get('allEnabled').updateValueAndValidity({onlySelf: true}); |
|||
} |
|||
|
|||
setDisabledState(isDisabled: boolean): void { |
|||
this.disabled = isDisabled; |
|||
if (isDisabled) { |
|||
this.debugSettingsFormGroup.disable({emitEvent: false}); |
|||
} else { |
|||
this.debugSettingsFormGroup.enable({emitEvent: false}); |
|||
} |
|||
} |
|||
} |
|||
@ -0,0 +1,66 @@ |
|||
<!-- |
|||
|
|||
Copyright © 2016-2024 The Thingsboard Authors |
|||
|
|||
Licensed under the Apache License, Version 2.0 (the "License"); |
|||
you may not use this file except in compliance with the License. |
|||
You may obtain a copy of the License at |
|||
|
|||
http://www.apache.org/licenses/LICENSE-2.0 |
|||
|
|||
Unless required by applicable law or agreed to in writing, software |
|||
distributed under the License is distributed on an "AS IS" BASIS, |
|||
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
See the License for the specific language governing permissions and |
|||
limitations under the License. |
|||
|
|||
--> |
|||
<div class="flex max-w-sm flex-col gap-3 p-2"> |
|||
<div class="tb-form-panel-title" translate>debug-config.label</div> |
|||
<div class="hint-container"> |
|||
<div class="tb-form-hint tb-primary-fill tb-flex center"> |
|||
@if (debugLimitsConfiguration) { |
|||
{{ 'debug-config.hint.main-limited' | translate: { msg: maxMessagesCount, sec: maxTimeFrameSec } }} |
|||
} @else { |
|||
{{ 'debug-config.hint.main' | translate }} |
|||
} |
|||
</div> |
|||
</div> |
|||
<div class="flex flex-col gap-3"> |
|||
<mat-slide-toggle class="mat-slide" [formControl]="onFailuresControl"> |
|||
<div tb-hint-tooltip-icon="{{ 'debug-config.hint.on-failure' | translate }}"> |
|||
{{ 'debug-config.on-failure' | translate }} |
|||
</div> |
|||
</mat-slide-toggle> |
|||
<div class="align-center flex justify-between"> |
|||
<mat-slide-toggle class="mat-slide" [formControl]="debugAllControl"> |
|||
<div tb-hint-tooltip-icon="{{ 'debug-config.hint.all-messages' | translate }}"> |
|||
{{ 'debug-config.all-messages' | translate: { time: (isDebugAllActive$ | async) && !allEnabled ? (allEnabledUntil | durationLeft) : ('debug-config.min' | translate: { number: maxDebugModeDurationMinutes }) } }} |
|||
</div> |
|||
</mat-slide-toggle> |
|||
<button mat-icon-button *ngIf="(isDebugAllActive$ | async) && !allEnabled" |
|||
class="tb-mat-20" |
|||
matTooltip="{{ 'action.reset' | translate }}" |
|||
matTooltipPosition="above" |
|||
color="primary" |
|||
(click)="onReset()"> |
|||
<mat-icon class="material-icons">refresh</mat-icon> |
|||
</button> |
|||
</div> |
|||
</div> |
|||
<div class="flex justify-end"> |
|||
<button mat-button |
|||
color="primary" |
|||
type="button" |
|||
(click)="onCancel()"> |
|||
{{ 'action.cancel' | translate }} |
|||
</button> |
|||
<button mat-raised-button |
|||
color="primary" |
|||
type="button" |
|||
(click)="onApply()"> |
|||
{{ 'action.apply' | translate }} |
|||
</button> |
|||
</div> |
|||
</div> |
|||
|
|||
@ -0,0 +1,134 @@ |
|||
///
|
|||
/// Copyright © 2016-2024 The Thingsboard Authors
|
|||
///
|
|||
/// Licensed under the Apache License, Version 2.0 (the "License");
|
|||
/// you may not use this file except in compliance with the License.
|
|||
/// You may obtain a copy of the License at
|
|||
///
|
|||
/// http://www.apache.org/licenses/LICENSE-2.0
|
|||
///
|
|||
/// Unless required by applicable law or agreed to in writing, software
|
|||
/// distributed under the License is distributed on an "AS IS" BASIS,
|
|||
/// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
|||
/// See the License for the specific language governing permissions and
|
|||
/// limitations under the License.
|
|||
///
|
|||
|
|||
import { |
|||
booleanAttribute, |
|||
ChangeDetectionStrategy, |
|||
ChangeDetectorRef, |
|||
Component, |
|||
EventEmitter, |
|||
Input, |
|||
OnInit |
|||
} from '@angular/core'; |
|||
import { PageComponent } from '@shared/components/page.component'; |
|||
import { TbPopoverComponent } from '@shared/components/popover.component'; |
|||
import { FormBuilder } from '@angular/forms'; |
|||
import { CommonModule } from '@angular/common'; |
|||
import { SharedModule } from '@shared/shared.module'; |
|||
import { SECOND } from '@shared/models/time/time.models'; |
|||
import { DurationLeftPipe } from '@shared/pipe/duration-left.pipe'; |
|||
import { of, shareReplay, timer } from 'rxjs'; |
|||
import { takeUntilDestroyed } from '@angular/core/rxjs-interop'; |
|||
import { DebugSettings } from '@shared/models/entity.models'; |
|||
import { distinctUntilChanged, map, startWith, switchMap, takeWhile } from 'rxjs/operators'; |
|||
|
|||
@Component({ |
|||
selector: 'tb-debug-settings-panel', |
|||
templateUrl: './debug-settings-panel.component.html', |
|||
standalone: true, |
|||
imports: [ |
|||
SharedModule, |
|||
CommonModule, |
|||
DurationLeftPipe |
|||
], |
|||
changeDetection: ChangeDetectionStrategy.OnPush |
|||
}) |
|||
export class DebugSettingsPanelComponent extends PageComponent implements OnInit { |
|||
|
|||
@Input() popover: TbPopoverComponent<DebugSettingsPanelComponent>; |
|||
@Input({ transform: booleanAttribute }) failuresEnabled = false; |
|||
@Input({ transform: booleanAttribute }) allEnabled = false; |
|||
@Input() allEnabledUntil = 0; |
|||
@Input() maxDebugModeDurationMinutes: number; |
|||
@Input() debugLimitsConfiguration: string; |
|||
|
|||
onFailuresControl = this.fb.control(false); |
|||
debugAllControl = this.fb.control(false); |
|||
|
|||
maxMessagesCount: string; |
|||
maxTimeFrameSec: string; |
|||
initialAllEnabled: boolean; |
|||
|
|||
isDebugAllActive$ = this.debugAllControl.valueChanges.pipe( |
|||
startWith(this.debugAllControl.value), |
|||
switchMap(value => { |
|||
if (value) { |
|||
return of(true); |
|||
} else { |
|||
return timer(0, SECOND).pipe( |
|||
map(() => this.allEnabledUntil > new Date().getTime()), |
|||
takeWhile(value => value, true) |
|||
); |
|||
} |
|||
}), |
|||
takeUntilDestroyed(), |
|||
shareReplay(1), |
|||
); |
|||
|
|||
onSettingsApplied = new EventEmitter<DebugSettings>(); |
|||
|
|||
constructor(private fb: FormBuilder, |
|||
private cd: ChangeDetectorRef) { |
|||
super(); |
|||
|
|||
this.debugAllControl.valueChanges.pipe( |
|||
takeUntilDestroyed() |
|||
).subscribe(value => { |
|||
this.allEnabled = value; |
|||
this.cd.markForCheck(); |
|||
}); |
|||
|
|||
this.isDebugAllActive$.pipe( |
|||
distinctUntilChanged(), |
|||
takeUntilDestroyed() |
|||
).subscribe(isDebugOn => this.debugAllControl.patchValue(isDebugOn, {emitEvent: false})) |
|||
} |
|||
|
|||
ngOnInit(): void { |
|||
this.maxMessagesCount = this.debugLimitsConfiguration?.split(':')[0]; |
|||
this.maxTimeFrameSec = this.debugLimitsConfiguration?.split(':')[1]; |
|||
this.onFailuresControl.patchValue(this.failuresEnabled); |
|||
this.debugAllControl.patchValue(this.allEnabled); |
|||
this.initialAllEnabled = this.allEnabled || this.allEnabledUntil > new Date().getTime(); |
|||
} |
|||
|
|||
onCancel(): void { |
|||
this.popover?.hide(); |
|||
} |
|||
|
|||
onApply(): void { |
|||
const isDebugAllChanged = this.initialAllEnabled !== this.debugAllControl.value || this.initialAllEnabled !== this.allEnabledUntil > new Date().getTime(); |
|||
if (isDebugAllChanged) { |
|||
this.onSettingsApplied.emit({ |
|||
allEnabled: this.allEnabled, |
|||
failuresEnabled: this.onFailuresControl.value, |
|||
allEnabledUntil: 0, |
|||
}); |
|||
} else { |
|||
this.onSettingsApplied.emit({ |
|||
allEnabled: false, |
|||
failuresEnabled: this.onFailuresControl.value, |
|||
allEnabledUntil: this.allEnabledUntil, |
|||
}); |
|||
} |
|||
} |
|||
|
|||
onReset(): void { |
|||
this.debugAllControl.patchValue(true); |
|||
this.allEnabledUntil = 0; |
|||
this.cd.markForCheck(); |
|||
} |
|||
} |
|||
@ -0,0 +1,35 @@ |
|||
///
|
|||
/// Copyright © 2016-2024 The Thingsboard Authors
|
|||
///
|
|||
/// Licensed under the Apache License, Version 2.0 (the "License");
|
|||
/// you may not use this file except in compliance with the License.
|
|||
/// You may obtain a copy of the License at
|
|||
///
|
|||
/// http://www.apache.org/licenses/LICENSE-2.0
|
|||
///
|
|||
/// Unless required by applicable law or agreed to in writing, software
|
|||
/// distributed under the License is distributed on an "AS IS" BASIS,
|
|||
/// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
|||
/// See the License for the specific language governing permissions and
|
|||
/// limitations under the License.
|
|||
///
|
|||
|
|||
import { Pipe, PipeTransform } from '@angular/core'; |
|||
import { TranslateService } from '@ngx-translate/core'; |
|||
import { MillisecondsToTimeStringPipe } from './milliseconds-to-time-string.pipe'; |
|||
|
|||
@Pipe({ |
|||
name: 'durationLeft', |
|||
pure: false, |
|||
standalone: true, |
|||
}) |
|||
export class DurationLeftPipe implements PipeTransform { |
|||
|
|||
constructor(private translate: TranslateService, private millisecondsToTimeString: MillisecondsToTimeStringPipe) { |
|||
} |
|||
|
|||
transform(untilTimestamp: number, shortFormat = true, onlyFirstDigit = true): string { |
|||
const time = this.millisecondsToTimeString.transform((untilTimestamp - new Date().getTime()), shortFormat, onlyFirstDigit) ?? 0; |
|||
return this.translate.instant('common.time-left', { time }); |
|||
} |
|||
} |
|||
Loading…
Reference in new issue