|
|
|
@ -16,20 +16,20 @@ |
|
|
|
package org.thingsboard.rule.engine.metadata; |
|
|
|
|
|
|
|
import com.google.common.util.concurrent.Futures; |
|
|
|
import lombok.Data; |
|
|
|
import lombok.RequiredArgsConstructor; |
|
|
|
import org.assertj.core.api.Assertions; |
|
|
|
import lombok.extern.slf4j.Slf4j; |
|
|
|
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.NullAndEmptySource; |
|
|
|
import org.junit.jupiter.params.provider.ValueSource; |
|
|
|
import org.mockito.ArgumentCaptor; |
|
|
|
import org.mockito.ArgumentMatcher; |
|
|
|
import org.mockito.Mock; |
|
|
|
import org.mockito.Spy; |
|
|
|
import org.mockito.junit.jupiter.MockitoExtension; |
|
|
|
import org.thingsboard.common.util.AbstractListeningExecutor; |
|
|
|
import org.thingsboard.common.util.JacksonUtil; |
|
|
|
import org.thingsboard.common.util.ListeningExecutor; |
|
|
|
import org.thingsboard.rule.engine.AbstractRuleNodeUpgradeTest; |
|
|
|
@ -54,34 +54,48 @@ import org.thingsboard.server.common.msg.TbMsgMetaData; |
|
|
|
import org.thingsboard.server.dao.timeseries.TimeseriesService; |
|
|
|
|
|
|
|
import java.util.List; |
|
|
|
import java.util.Optional; |
|
|
|
import java.util.UUID; |
|
|
|
import java.util.concurrent.CountDownLatch; |
|
|
|
import java.util.concurrent.TimeUnit; |
|
|
|
import java.util.function.BiConsumer; |
|
|
|
import java.util.stream.IntStream; |
|
|
|
import java.util.stream.Stream; |
|
|
|
|
|
|
|
import static org.assertj.core.api.Assertions.assertThat; |
|
|
|
import static org.awaitility.Awaitility.await; |
|
|
|
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; |
|
|
|
import static org.junit.jupiter.api.Assertions.assertEquals; |
|
|
|
import static org.junit.jupiter.api.Assertions.assertFalse; |
|
|
|
import static org.junit.jupiter.api.Assertions.assertInstanceOf; |
|
|
|
import static org.junit.jupiter.api.Assertions.assertThrows; |
|
|
|
import static org.junit.jupiter.api.Assertions.assertTrue; |
|
|
|
import static org.mockito.ArgumentMatchers.any; |
|
|
|
import static org.mockito.ArgumentMatchers.anyList; |
|
|
|
import static org.mockito.ArgumentMatchers.anySet; |
|
|
|
import static org.mockito.ArgumentMatchers.anyString; |
|
|
|
import static org.mockito.ArgumentMatchers.argThat; |
|
|
|
import static org.mockito.ArgumentMatchers.eq; |
|
|
|
import static org.mockito.BDDMockito.willAnswer; |
|
|
|
import static org.mockito.Mockito.mock; |
|
|
|
import static org.mockito.Mockito.never; |
|
|
|
import static org.mockito.Mockito.reset; |
|
|
|
import static org.mockito.Mockito.spy; |
|
|
|
import static org.mockito.Mockito.times; |
|
|
|
import static org.mockito.Mockito.verify; |
|
|
|
import static org.mockito.Mockito.verifyNoMoreInteractions; |
|
|
|
import static org.mockito.Mockito.when; |
|
|
|
|
|
|
|
@Slf4j |
|
|
|
@ExtendWith(MockitoExtension.class) |
|
|
|
public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest { |
|
|
|
|
|
|
|
private static final DeviceId DUMMY_DEVICE_ORIGINATOR = new DeviceId(UUID.randomUUID()); |
|
|
|
private static final TenantId TENANT_ID = new TenantId(UUID.randomUUID()); |
|
|
|
private static final ListeningExecutor DB_EXECUTOR = new TestDbCallbackExecutor(); |
|
|
|
private final DeviceId DUMMY_DEVICE_ORIGINATOR = new DeviceId(UUID.fromString("2ba3ded4-882b-40cf-999a-89da9ccd58f9")); |
|
|
|
private final TenantId TENANT_ID = TenantId.fromUUID(UUID.fromString("3842e740-0d89-43a9-8d52-ae44023847ba")); |
|
|
|
private final ListeningExecutor DB_EXECUTOR = new TestDbCallbackExecutor(); |
|
|
|
|
|
|
|
private static final int RULE_DISPATCHER_POOL_SIZE = 2; |
|
|
|
private static final int DB_CALLBACK_POOL_SIZE = 3; |
|
|
|
|
|
|
|
@Mock |
|
|
|
private TbContext ctxMock; |
|
|
|
@Mock |
|
|
|
@ -95,8 +109,6 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest { |
|
|
|
public void setUp() throws TbNodeException { |
|
|
|
config = new CalculateDeltaNodeConfiguration().defaultConfiguration(); |
|
|
|
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|
|
|
when(ctxMock.getTimeseriesService()).thenReturn(timeseriesServiceMock); |
|
|
|
|
|
|
|
node.init(ctxMock, nodeConfiguration); |
|
|
|
} |
|
|
|
|
|
|
|
@ -110,6 +122,49 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest { |
|
|
|
assertTrue(config.isTellFailureIfDeltaIsNegative()); |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
@ParameterizedTest |
|
|
|
@NullAndEmptySource |
|
|
|
@ValueSource(strings = {" "}) // blank value
|
|
|
|
public void givenInvalidInputKey_whenInitThenThrowException(String key) { |
|
|
|
config.setInputValueKey(key); |
|
|
|
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|
|
|
var exception = assertThrows(TbNodeException.class, () -> node.init(ctxMock, nodeConfiguration)); |
|
|
|
assertThat(exception).hasMessage("Input value key should be specified!"); |
|
|
|
assertThat(exception.isUnrecoverable()).isTrue(); |
|
|
|
} |
|
|
|
|
|
|
|
@ParameterizedTest |
|
|
|
@NullAndEmptySource |
|
|
|
@ValueSource(strings = {" "}) // blank value
|
|
|
|
public void givenInvalidOutputKey_whenInitThenThrowException(String key) { |
|
|
|
config.setOutputValueKey(key); |
|
|
|
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|
|
|
var exception = assertThrows(TbNodeException.class, () -> node.init(ctxMock, nodeConfiguration)); |
|
|
|
assertThat(exception).hasMessage("Output value key should be specified!"); |
|
|
|
assertThat(exception.isUnrecoverable()).isTrue(); |
|
|
|
} |
|
|
|
|
|
|
|
@ParameterizedTest |
|
|
|
@NullAndEmptySource |
|
|
|
@ValueSource(strings = {" "}) // blank value
|
|
|
|
public void givenInvalidPeriodKey_whenInitThenThrowException(String key) { |
|
|
|
config.setPeriodValueKey(key); |
|
|
|
config.setAddPeriodBetweenMsgs(true); |
|
|
|
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|
|
|
var exception = assertThrows(TbNodeException.class, () -> node.init(ctxMock, nodeConfiguration)); |
|
|
|
assertThat(exception).hasMessage("Period value key should be specified!"); |
|
|
|
assertThat(exception.isUnrecoverable()).isTrue(); |
|
|
|
} |
|
|
|
|
|
|
|
@Test |
|
|
|
public void givenInvalidPeriodKeyAndAddPeriodDisabled_whenInitThenNoExceptionThrown() { |
|
|
|
config.setPeriodValueKey(null); |
|
|
|
config.setAddPeriodBetweenMsgs(false); |
|
|
|
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|
|
|
assertDoesNotThrow(() -> node.init(ctxMock, nodeConfiguration)); |
|
|
|
} |
|
|
|
|
|
|
|
@Test |
|
|
|
public void givenInvalidMsgType_whenOnMsg_thenShouldTellNextOther() { |
|
|
|
// GIVEN
|
|
|
|
@ -120,7 +175,7 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest { |
|
|
|
node.onMsg(ctxMock, msg); |
|
|
|
|
|
|
|
// THEN
|
|
|
|
verify(ctxMock, times(1)).tellNext(eq(msg), eq(TbNodeConnectionType.OTHER)); |
|
|
|
verify(ctxMock).tellNext(eq(msg), eq(TbNodeConnectionType.OTHER)); |
|
|
|
verify(ctxMock, never()).tellSuccess(any()); |
|
|
|
verify(ctxMock, never()).tellFailure(any(), any()); |
|
|
|
} |
|
|
|
@ -134,7 +189,7 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest { |
|
|
|
node.onMsg(ctxMock, msg); |
|
|
|
|
|
|
|
// THEN
|
|
|
|
verify(ctxMock, times(1)).tellNext(eq(msg), eq(TbNodeConnectionType.OTHER)); |
|
|
|
verify(ctxMock).tellNext(eq(msg), eq(TbNodeConnectionType.OTHER)); |
|
|
|
verify(ctxMock, never()).tellSuccess(any()); |
|
|
|
verify(ctxMock, never()).tellFailure(any(), any()); |
|
|
|
} |
|
|
|
@ -149,7 +204,7 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest { |
|
|
|
node.onMsg(ctxMock, msg); |
|
|
|
|
|
|
|
// THEN
|
|
|
|
verify(ctxMock, times(1)).tellNext(eq(msg), eq(TbNodeConnectionType.OTHER)); |
|
|
|
verify(ctxMock).tellNext(eq(msg), eq(TbNodeConnectionType.OTHER)); |
|
|
|
verify(ctxMock, never()).tellSuccess(any()); |
|
|
|
verify(ctxMock, never()).tellFailure(any(), any()); |
|
|
|
} |
|
|
|
@ -175,7 +230,7 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest { |
|
|
|
// THEN
|
|
|
|
var actualMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|
|
|
|
|
|
|
verify(ctxMock, times(1)).tellSuccess(actualMsgCaptor.capture()); |
|
|
|
verify(ctxMock).tellSuccess(actualMsgCaptor.capture()); |
|
|
|
verify(ctxMock, never()).tellNext(any(), anyString()); |
|
|
|
verify(ctxMock, never()).tellNext(any(), anySet()); |
|
|
|
verify(ctxMock, never()).tellFailure(any(), any()); |
|
|
|
@ -205,7 +260,7 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest { |
|
|
|
// THEN
|
|
|
|
var actualMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|
|
|
|
|
|
|
verify(ctxMock, times(1)).tellSuccess(actualMsgCaptor.capture()); |
|
|
|
verify(ctxMock).tellSuccess(actualMsgCaptor.capture()); |
|
|
|
verify(ctxMock, never()).tellNext(any(), anyString()); |
|
|
|
verify(ctxMock, never()).tellNext(any(), anySet()); |
|
|
|
verify(ctxMock, never()).tellFailure(any(), any()); |
|
|
|
@ -235,7 +290,7 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest { |
|
|
|
// THEN
|
|
|
|
var actualMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|
|
|
|
|
|
|
verify(ctxMock, times(1)).tellSuccess(actualMsgCaptor.capture()); |
|
|
|
verify(ctxMock).tellSuccess(actualMsgCaptor.capture()); |
|
|
|
verify(ctxMock, never()).tellNext(any(), anyString()); |
|
|
|
verify(ctxMock, never()).tellNext(any(), anySet()); |
|
|
|
verify(ctxMock, never()).tellFailure(any(), any()); |
|
|
|
@ -256,7 +311,7 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest { |
|
|
|
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|
|
|
node.init(ctxMock, nodeConfiguration); |
|
|
|
|
|
|
|
mockFindLatest(new BasicTsKvEntry(1L, new DoubleDataEntry("temperature", 40.0))); |
|
|
|
mockFindLatestAsync(new BasicTsKvEntry(1L, new DoubleDataEntry("temperature", 40.0))); |
|
|
|
|
|
|
|
var msgData = "{\"temperature\": 42,\"airPressure\":123}"; |
|
|
|
var firstMsgMetaData = new TbMsgMetaData(); |
|
|
|
@ -269,7 +324,7 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest { |
|
|
|
// THEN
|
|
|
|
var actualMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|
|
|
|
|
|
|
verify(ctxMock, times(1)).tellSuccess(actualMsgCaptor.capture()); |
|
|
|
verify(ctxMock).tellSuccess(actualMsgCaptor.capture()); |
|
|
|
verify(ctxMock, never()).tellNext(any(), anyString()); |
|
|
|
verify(ctxMock, never()).tellNext(any(), anySet()); |
|
|
|
verify(ctxMock, never()).tellFailure(any(), any()); |
|
|
|
@ -283,6 +338,8 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest { |
|
|
|
reset(ctxMock); |
|
|
|
reset(timeseriesServiceMock); |
|
|
|
|
|
|
|
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR); |
|
|
|
|
|
|
|
var secondMsgMetaData = new TbMsgMetaData(); |
|
|
|
secondMsgMetaData.putValue("ts", String.valueOf(6L)); |
|
|
|
var secondMsg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DUMMY_DEVICE_ORIGINATOR, secondMsgMetaData, msgData); |
|
|
|
@ -294,7 +351,7 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest { |
|
|
|
actualMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|
|
|
|
|
|
|
verify(timeseriesServiceMock, never()).findLatest(any(), any(), anyList()); |
|
|
|
verify(ctxMock, times(1)).tellSuccess(actualMsgCaptor.capture()); |
|
|
|
verify(ctxMock).tellSuccess(actualMsgCaptor.capture()); |
|
|
|
verify(ctxMock, never()).tellNext(any(), anyString()); |
|
|
|
verify(ctxMock, never()).tellNext(any(), anySet()); |
|
|
|
verify(ctxMock, never()).tellFailure(any(), any()); |
|
|
|
@ -324,7 +381,7 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest { |
|
|
|
// THEN
|
|
|
|
var actualMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|
|
|
|
|
|
|
verify(ctxMock, times(1)).tellSuccess(actualMsgCaptor.capture()); |
|
|
|
verify(ctxMock).tellSuccess(actualMsgCaptor.capture()); |
|
|
|
verify(ctxMock, never()).tellNext(any(), anyString()); |
|
|
|
verify(ctxMock, never()).tellNext(any(), anySet()); |
|
|
|
verify(ctxMock, never()).tellFailure(any(), any()); |
|
|
|
@ -341,7 +398,7 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest { |
|
|
|
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|
|
|
node.init(ctxMock, nodeConfiguration); |
|
|
|
|
|
|
|
mockFindLatest(new BasicTsKvEntry(System.currentTimeMillis(), new LongDataEntry("pulseCounter", 200L))); |
|
|
|
mockFindLatestAsync(new BasicTsKvEntry(System.currentTimeMillis(), new LongDataEntry("pulseCounter", 200L))); |
|
|
|
|
|
|
|
var msgData = "{\"pulseCounter\":\"123\"}"; |
|
|
|
var msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DUMMY_DEVICE_ORIGINATOR, TbMsgMetaData.EMPTY, msgData); |
|
|
|
@ -353,7 +410,7 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest { |
|
|
|
var actualMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|
|
|
var actualExceptionCaptor = ArgumentCaptor.forClass(Exception.class); |
|
|
|
|
|
|
|
verify(ctxMock, times(1)).tellFailure(actualMsgCaptor.capture(), actualExceptionCaptor.capture()); |
|
|
|
verify(ctxMock).tellFailure(actualMsgCaptor.capture(), actualExceptionCaptor.capture()); |
|
|
|
verify(ctxMock, never()).tellSuccess(any()); |
|
|
|
verify(ctxMock, never()).tellNext(any(), anyString()); |
|
|
|
verify(ctxMock, never()).tellNext(any(), anySet()); |
|
|
|
@ -373,7 +430,7 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest { |
|
|
|
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|
|
|
node.init(ctxMock, nodeConfiguration); |
|
|
|
|
|
|
|
mockFindLatest(new BasicTsKvEntry(System.currentTimeMillis(), new LongDataEntry("pulseCounter", 200L))); |
|
|
|
mockFindLatestAsync(new BasicTsKvEntry(System.currentTimeMillis(), new LongDataEntry("pulseCounter", 200L))); |
|
|
|
|
|
|
|
var msgData = "{\"pulseCounter\":\"123\"}"; |
|
|
|
var msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DUMMY_DEVICE_ORIGINATOR, TbMsgMetaData.EMPTY, msgData); |
|
|
|
@ -384,7 +441,7 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest { |
|
|
|
// THEN
|
|
|
|
var actualMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|
|
|
|
|
|
|
verify(ctxMock, times(1)).tellSuccess(actualMsgCaptor.capture()); |
|
|
|
verify(ctxMock).tellSuccess(actualMsgCaptor.capture()); |
|
|
|
verify(ctxMock, never()).tellFailure(any(), any()); |
|
|
|
verify(ctxMock, never()).tellNext(any(), anyString()); |
|
|
|
verify(ctxMock, never()).tellNext(any(), anySet()); |
|
|
|
@ -396,13 +453,23 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest { |
|
|
|
@Test |
|
|
|
public void givenInvalidStringValue_whenOnMsg_thenException() { |
|
|
|
// GIVEN
|
|
|
|
mockFindLatest(new BasicTsKvEntry(System.currentTimeMillis(), new StringDataEntry("pulseCounter", "high"))); |
|
|
|
mockFindLatestAsync(new BasicTsKvEntry(System.currentTimeMillis(), new StringDataEntry("pulseCounter", "high"))); |
|
|
|
|
|
|
|
var msgData = "{\"pulseCounter\":\"123\"}"; |
|
|
|
var msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DUMMY_DEVICE_ORIGINATOR, TbMsgMetaData.EMPTY, msgData); |
|
|
|
|
|
|
|
// WHEN-THEN
|
|
|
|
Assertions.assertThatThrownBy(() -> node.onMsg(ctxMock, msg)) |
|
|
|
// WHEN
|
|
|
|
node.onMsg(ctxMock, msg); |
|
|
|
|
|
|
|
// THEN
|
|
|
|
ArgumentCaptor<Throwable> throwableCaptor = ArgumentCaptor.forClass(Throwable.class); |
|
|
|
|
|
|
|
verify(ctxMock).tellFailure(eq(msg), throwableCaptor.capture()); |
|
|
|
verify(ctxMock, never()).tellSuccess(any()); |
|
|
|
verify(ctxMock, never()).tellNext(any(), anyString()); |
|
|
|
verify(ctxMock, never()).tellNext(any(), anySet()); |
|
|
|
|
|
|
|
assertThat(throwableCaptor.getValue()) |
|
|
|
.isInstanceOf(IllegalArgumentException.class) |
|
|
|
.hasMessage("Calculation failed. Unable to parse value [high] of telemetry [pulseCounter] to Double"); |
|
|
|
} |
|
|
|
@ -410,13 +477,23 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest { |
|
|
|
@Test |
|
|
|
public void givenBooleanValue_whenOnMsg_thenException() { |
|
|
|
// GIVEN
|
|
|
|
mockFindLatest(new BasicTsKvEntry(System.currentTimeMillis(), new BooleanDataEntry("pulseCounter", false))); |
|
|
|
mockFindLatestAsync(new BasicTsKvEntry(System.currentTimeMillis(), new BooleanDataEntry("pulseCounter", false))); |
|
|
|
|
|
|
|
var msgData = "{\"pulseCounter\":true}"; |
|
|
|
var msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DUMMY_DEVICE_ORIGINATOR, TbMsgMetaData.EMPTY, msgData); |
|
|
|
|
|
|
|
// WHEN-THEN
|
|
|
|
Assertions.assertThatThrownBy(() -> node.onMsg(ctxMock, msg)) |
|
|
|
// WHEN
|
|
|
|
node.onMsg(ctxMock, msg); |
|
|
|
|
|
|
|
// THEN
|
|
|
|
ArgumentCaptor<Throwable> throwableCaptor = ArgumentCaptor.forClass(Throwable.class); |
|
|
|
|
|
|
|
verify(ctxMock).tellFailure(eq(msg), throwableCaptor.capture()); |
|
|
|
verify(ctxMock, never()).tellSuccess(any()); |
|
|
|
verify(ctxMock, never()).tellNext(any(), anyString()); |
|
|
|
verify(ctxMock, never()).tellNext(any(), anySet()); |
|
|
|
|
|
|
|
assertThat(throwableCaptor.getValue()) |
|
|
|
.isInstanceOf(IllegalArgumentException.class) |
|
|
|
.hasMessage("Calculation failed. Boolean values are not supported!"); |
|
|
|
} |
|
|
|
@ -424,38 +501,104 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest { |
|
|
|
@Test |
|
|
|
public void givenJsonValue_whenOnMsg_thenException() { |
|
|
|
// GIVEN
|
|
|
|
mockFindLatest(new BasicTsKvEntry(System.currentTimeMillis(), new JsonDataEntry("pulseCounter", "{\"isActive\":false}"))); |
|
|
|
mockFindLatestAsync(new BasicTsKvEntry(System.currentTimeMillis(), new JsonDataEntry("pulseCounter", "{\"isActive\":false}"))); |
|
|
|
|
|
|
|
var msgData = "{\"pulseCounter\":{\"isActive\":true}}"; |
|
|
|
var msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DUMMY_DEVICE_ORIGINATOR, TbMsgMetaData.EMPTY, msgData); |
|
|
|
|
|
|
|
// WHEN-THEN
|
|
|
|
Assertions.assertThatThrownBy(() -> node.onMsg(ctxMock, msg)) |
|
|
|
// WHEN
|
|
|
|
node.onMsg(ctxMock, msg); |
|
|
|
|
|
|
|
// THEN
|
|
|
|
ArgumentCaptor<Throwable> throwableCaptor = ArgumentCaptor.forClass(Throwable.class); |
|
|
|
|
|
|
|
verify(ctxMock).tellFailure(eq(msg), throwableCaptor.capture()); |
|
|
|
verify(ctxMock, never()).tellSuccess(any()); |
|
|
|
verify(ctxMock, never()).tellNext(any(), anyString()); |
|
|
|
verify(ctxMock, never()).tellNext(any(), anySet()); |
|
|
|
|
|
|
|
assertThat(throwableCaptor.getValue()) |
|
|
|
.isInstanceOf(IllegalArgumentException.class) |
|
|
|
.hasMessage("Calculation failed. JSON values are not supported!"); |
|
|
|
} |
|
|
|
|
|
|
|
@Test |
|
|
|
public void givenConcurrentAccess_whenOnMsg_thenGetFromDBInvokedOnce() throws TbNodeException, InterruptedException { |
|
|
|
DBCallbackExecutor dbCallbackExecutor = new DBCallbackExecutor(); |
|
|
|
dbCallbackExecutor.init(); |
|
|
|
|
|
|
|
RuleDispatcherExecutor ruleEngineDispatcherExecutor = new RuleDispatcherExecutor(); |
|
|
|
ruleEngineDispatcherExecutor.init(); |
|
|
|
|
|
|
|
assertThat(RULE_DISPATCHER_POOL_SIZE).as("dispatcher pool size have to be > 1").isGreaterThan(1); |
|
|
|
|
|
|
|
final TbContext ctx = mock(TbContext.class); |
|
|
|
final TimeseriesService timeseriesService = mock(TimeseriesService.class); |
|
|
|
|
|
|
|
when(ctx.getTimeseriesService()).thenReturn(timeseriesService); |
|
|
|
when(ctx.getDbCallbackExecutor()).thenReturn(dbCallbackExecutor); |
|
|
|
when(timeseriesService.findLatest(any(), any(), anyString())).thenReturn(Futures.immediateFuture(Optional.empty())); |
|
|
|
|
|
|
|
final CalculateDeltaNodeConfiguration config = new CalculateDeltaNodeConfiguration().defaultConfiguration(); |
|
|
|
final TbNodeConfiguration nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|
|
|
final CalculateDeltaNode node = spy(CalculateDeltaNode.class); |
|
|
|
|
|
|
|
node.init(ctx, nodeConfiguration); |
|
|
|
|
|
|
|
List<TbMsg> tbMsgList = IntStream.range(0, RULE_DISPATCHER_POOL_SIZE * 2).mapToObj(x -> { |
|
|
|
var msgData = "{\"pulseCounter\":" + 2 + "}"; |
|
|
|
return TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DUMMY_DEVICE_ORIGINATOR, TbMsgMetaData.EMPTY, msgData); |
|
|
|
}).toList(); |
|
|
|
|
|
|
|
CountDownLatch processingLatch = new CountDownLatch(tbMsgList.size()); |
|
|
|
|
|
|
|
willAnswer(invocation -> { |
|
|
|
processingLatch.countDown(); |
|
|
|
return invocation.callRealMethod(); |
|
|
|
}).given(node).processMsgAsync(any(), any()); |
|
|
|
|
|
|
|
tbMsgList.forEach(msg -> ruleEngineDispatcherExecutor.executeAsync(() -> node.onMsg(ctx, msg))); |
|
|
|
|
|
|
|
assertThat(processingLatch.await(5, TimeUnit.SECONDS)).as("await on processingLatch").isTrue(); |
|
|
|
|
|
|
|
verify(timeseriesService).findLatest(any(), any(), anyString()); |
|
|
|
await().atMost(5, TimeUnit.SECONDS).untilAsserted(() -> verify(ctx, times(tbMsgList.size())).tellSuccess(any())); |
|
|
|
} |
|
|
|
|
|
|
|
private static class RuleDispatcherExecutor extends AbstractListeningExecutor { |
|
|
|
@Override |
|
|
|
protected int getThreadPollSize() { |
|
|
|
return RULE_DISPATCHER_POOL_SIZE; |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private static class DBCallbackExecutor extends AbstractListeningExecutor { |
|
|
|
@Override |
|
|
|
protected int getThreadPollSize() { |
|
|
|
return DB_CALLBACK_POOL_SIZE; |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
@ParameterizedTest |
|
|
|
@MethodSource("CalculateDeltaTestConfig") |
|
|
|
public void givenCalculateDeltaConfig_whenOnMsg_thenVerify(CalculateDeltaTestConfig testConfig) throws TbNodeException { |
|
|
|
// GIVEN
|
|
|
|
config.setTellFailureIfDeltaIsNegative(testConfig.isTellFailureIfDeltaIsNegative()); |
|
|
|
config.setExcludeZeroDeltas(testConfig.isExcludeZeroDeltas()); |
|
|
|
config.setTellFailureIfDeltaIsNegative(testConfig.tellFailureIfDeltaIsNegative()); |
|
|
|
config.setExcludeZeroDeltas(testConfig.excludeZeroDeltas()); |
|
|
|
config.setInputValueKey("temperature"); |
|
|
|
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|
|
|
node.init(ctxMock, nodeConfiguration); |
|
|
|
|
|
|
|
mockFindLatest(new BasicTsKvEntry(1L, new DoubleDataEntry("temperature", testConfig.getPrevValue()))); |
|
|
|
mockFindLatestAsync(new BasicTsKvEntry(1L, new DoubleDataEntry("temperature", testConfig.prevValue()))); |
|
|
|
|
|
|
|
var msgData = "{\"temperature\":" + testConfig.getCurrentValue() + ",\"airPressure\":123}"; |
|
|
|
var msgData = "{\"temperature\":" + testConfig.currentValue() + ",\"airPressure\":123}"; |
|
|
|
var msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DUMMY_DEVICE_ORIGINATOR, TbMsgMetaData.EMPTY, msgData); |
|
|
|
|
|
|
|
// WHEN
|
|
|
|
|
|
|
|
node.onMsg(ctxMock, msg); |
|
|
|
|
|
|
|
// THEN
|
|
|
|
testConfig.getVerificationMethod().accept(ctxMock, msg); |
|
|
|
testConfig.verificationMethod().accept(ctxMock, msg); |
|
|
|
} |
|
|
|
|
|
|
|
private static Stream<CalculateDeltaTestConfig> CalculateDeltaTestConfig() { |
|
|
|
@ -510,47 +653,18 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest { |
|
|
|
); |
|
|
|
} |
|
|
|
|
|
|
|
@Data |
|
|
|
@RequiredArgsConstructor |
|
|
|
private static class CalculateDeltaTestConfig { |
|
|
|
private final boolean tellFailureIfDeltaIsNegative; |
|
|
|
private final boolean excludeZeroDeltas; |
|
|
|
private final double prevValue; |
|
|
|
private final double currentValue; |
|
|
|
private final BiConsumer<TbContext, TbMsg> verificationMethod; |
|
|
|
} |
|
|
|
|
|
|
|
private void mockFindLatest(TsKvEntry tsKvEntry) { |
|
|
|
when(ctxMock.getTenantId()).thenReturn(TENANT_ID); |
|
|
|
when(timeseriesServiceMock.findLatestSync( |
|
|
|
eq(TENANT_ID), eq(DUMMY_DEVICE_ORIGINATOR), argThat(new ListMatcher<>(List.of(tsKvEntry.getKey()))) |
|
|
|
)).thenReturn(List.of(tsKvEntry)); |
|
|
|
private record CalculateDeltaTestConfig(boolean tellFailureIfDeltaIsNegative, boolean excludeZeroDeltas, |
|
|
|
double prevValue, double currentValue, |
|
|
|
BiConsumer<TbContext, TbMsg> verificationMethod) { |
|
|
|
} |
|
|
|
|
|
|
|
private void mockFindLatestAsync(TsKvEntry tsKvEntry) { |
|
|
|
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR); |
|
|
|
when(ctxMock.getTenantId()).thenReturn(TENANT_ID); |
|
|
|
when(ctxMock.getTimeseriesService()).thenReturn(timeseriesServiceMock); |
|
|
|
when(timeseriesServiceMock.findLatest( |
|
|
|
eq(TENANT_ID), eq(DUMMY_DEVICE_ORIGINATOR), argThat(new ListMatcher<>(List.of(tsKvEntry.getKey()))) |
|
|
|
)).thenReturn(Futures.immediateFuture(List.of(tsKvEntry))); |
|
|
|
} |
|
|
|
|
|
|
|
@RequiredArgsConstructor |
|
|
|
private static class ListMatcher<T> implements ArgumentMatcher<List<T>> { |
|
|
|
|
|
|
|
private final List<T> expectedList; |
|
|
|
|
|
|
|
@Override |
|
|
|
public boolean matches(List<T> actualList) { |
|
|
|
if (actualList == expectedList) { |
|
|
|
return true; |
|
|
|
} |
|
|
|
if (actualList.size() != expectedList.size()) { |
|
|
|
return false; |
|
|
|
} |
|
|
|
return actualList.containsAll(expectedList); |
|
|
|
} |
|
|
|
|
|
|
|
eq(TENANT_ID), eq(DUMMY_DEVICE_ORIGINATOR), eq(tsKvEntry.getKey()) |
|
|
|
)).thenReturn(Futures.immediateFuture(Optional.of(tsKvEntry))); |
|
|
|
} |
|
|
|
|
|
|
|
private static Stream<Arguments> givenFromVersionAndConfig_whenUpgrade_thenVerifyHasChangesAndConfig() { |
|
|
|
|