Browse Source

added verification on init

pull/11140/head
IrynaMatveieva 2 years ago
parent
commit
a34088fcfb
  1. 66
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/delay/TbMsgDelayNode.java
  2. 249
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/delay/TbMsgDelayNodeTest.java

66
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/delay/TbMsgDelayNode.java

@ -31,10 +31,10 @@ import org.thingsboard.server.common.data.util.TbPair;
import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgMetaData; import org.thingsboard.server.common.msg.TbMsgMetaData;
import java.util.HashMap;
import java.util.List; import java.util.List;
import java.util.Map;
import java.util.UUID; import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
@Slf4j @Slf4j
@ -54,13 +54,17 @@ import java.util.concurrent.TimeUnit;
) )
public class TbMsgDelayNode implements TbNode { public class TbMsgDelayNode implements TbNode {
private final List<TimeUnit> supportedTimeUnits = List.of(TimeUnit.SECONDS, TimeUnit.MINUTES, TimeUnit.HOURS);
private final String supportedTimeUnitsStr = String.join(",", TimeUnit.SECONDS.name(), TimeUnit.MINUTES.name(), TimeUnit.HOURS.name());
private TbMsgDelayNodeConfiguration config; private TbMsgDelayNodeConfiguration config;
private Map<UUID, TbMsg> pendingMsgs; private ConcurrentMap<UUID, TbMsg> pendingMsgs;
@Override @Override
public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException { public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException {
this.config = TbNodeUtils.convert(configuration, TbMsgDelayNodeConfiguration.class); this.config = TbNodeUtils.convert(configuration, TbMsgDelayNodeConfiguration.class);
this.pendingMsgs = new HashMap<>(); validateConfig();
this.pendingMsgs = new ConcurrentHashMap<>();
} }
@Override @Override
@ -71,7 +75,7 @@ public class TbMsgDelayNode implements TbNode {
ctx.enqueueForTellNext( ctx.enqueueForTellNext(
TbMsg.newMsg( TbMsg.newMsg(
pendingMsg.getQueueName(), pendingMsg.getQueueName(),
pendingMsg.getType(), pendingMsg.getInternalType(),
pendingMsg.getOriginator(), pendingMsg.getOriginator(),
pendingMsg.getCustomerId(), pendingMsg.getCustomerId(),
pendingMsg.getMetaData(), pendingMsg.getMetaData(),
@ -97,12 +101,28 @@ public class TbMsgDelayNode implements TbNode {
String periodPattern = TbNodeUtils.processPattern(config.getPeriod(), msg); String periodPattern = TbNodeUtils.processPattern(config.getPeriod(), msg);
try { try {
TimeUnit timeUnit = TimeUnit.valueOf(timeUnitPattern.toUpperCase()); TimeUnit timeUnit = TimeUnit.valueOf(timeUnitPattern.toUpperCase());
if (!supportedTimeUnits.contains(timeUnit)) {
throw new RuntimeException("Time unit '" + timeUnit + "' is not supported! " +
"Only " + supportedTimeUnitsStr + " are supported.");
}
int period = Integer.parseInt(periodPattern); int period = Integer.parseInt(periodPattern);
return timeUnit.toMillis(period); return timeUnit.toMillis(period);
} catch (NumberFormatException e) { } catch (NumberFormatException e) {
throw new RuntimeException("Can't parse period value : " + periodPattern); throw new NumberFormatException("Can't parse period value : " + periodPattern);
} catch (IllegalArgumentException e) { } catch (IllegalArgumentException e) {
throw new RuntimeException("Invalid value for period time unit : " + timeUnitPattern); throw new IllegalArgumentException("Invalid value for period time unit : " + timeUnitPattern);
}
}
private void validateConfig() throws TbNodeException {
if (config.getMaxPendingMsgs() < 1 || config.getMaxPendingMsgs() > 100000) {
throw new TbNodeException("Maximum pending messages should be in a range from 1 to 100000.", true);
}
if (config.getPeriod() == null) {
throw new TbNodeException("Period should be specified.", true);
}
if (config.getTimeUnit() == null) {
throw new TbNodeException("Time unit should be specified.", true);
} }
} }
@ -121,26 +141,30 @@ public class TbMsgDelayNode implements TbNode {
var useMetadataPeriodInSecondsPatterns = "useMetadataPeriodInSecondsPatterns"; var useMetadataPeriodInSecondsPatterns = "useMetadataPeriodInSecondsPatterns";
var period = "period"; var period = "period";
if (oldConfiguration.has(useMetadataPeriodInSecondsPatterns)) { if (oldConfiguration.has(useMetadataPeriodInSecondsPatterns)) {
var isUsedPattern = oldConfiguration.get(useMetadataPeriodInSecondsPatterns).asBoolean(); var isUsedPattern = oldConfiguration.get(useMetadataPeriodInSecondsPatterns).booleanValue();
if (isUsedPattern) { if (isUsedPattern) {
if (!oldConfiguration.has(periodInSecondsPattern)) { if (!oldConfiguration.has(periodInSecondsPattern)) {
throw new TbNodeException("Property to update: '" + periodInSecondsPattern + "' does not exist in configuration."); throw new TbNodeException("Property to update: '" + periodInSecondsPattern + "' does not exist in configuration.");
} }
((ObjectNode) oldConfiguration).set(period, oldConfiguration.get(periodInSecondsPattern)); ((ObjectNode) oldConfiguration).set(period, oldConfiguration.get(periodInSecondsPattern));
} else { } else {
if (!oldConfiguration.has(periodInSeconds)) { if (!oldConfiguration.has(periodInSeconds)) {
throw new TbNodeException("Property to update: '" + periodInSeconds + "' does not exist in configuration."); throw new TbNodeException("Property to update: '" + periodInSeconds + "' does not exist in configuration.");
} }
((ObjectNode) oldConfiguration).set(period, oldConfiguration.get(periodInSeconds)); ((ObjectNode) oldConfiguration).put(period, oldConfiguration.get(periodInSeconds).asText());
} }
((ObjectNode) oldConfiguration).remove(List.of(periodInSeconds, periodInSecondsPattern, useMetadataPeriodInSecondsPatterns)); hasChanges = true;
hasChanges = true; }
if (!oldConfiguration.has(period)) {
((ObjectNode) oldConfiguration).put(period, "60");
hasChanges = true;
} }
var timeUnit = "timeUnit"; var timeUnit = "timeUnit";
if (!oldConfiguration.has(timeUnit)) { if (!oldConfiguration.has(timeUnit)) {
((ObjectNode) oldConfiguration).put(timeUnit, TimeUnit.SECONDS.name()); ((ObjectNode) oldConfiguration).put(timeUnit, TimeUnit.SECONDS.name());
hasChanges = true; hasChanges = true;
} }
((ObjectNode) oldConfiguration).remove(List.of(periodInSeconds, periodInSecondsPattern, useMetadataPeriodInSecondsPatterns));
break; break;
default: default:
break; break;

249
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/delay/TbMsgDelayNodeTest.java

@ -20,7 +20,9 @@ import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith; import org.junit.jupiter.api.extension.ExtendWith;
import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.Arguments; import org.junit.jupiter.params.provider.Arguments;
import org.junit.jupiter.params.provider.EnumSource;
import org.junit.jupiter.params.provider.MethodSource; import org.junit.jupiter.params.provider.MethodSource;
import org.junit.jupiter.params.provider.ValueSource;
import org.mockito.ArgumentCaptor; import org.mockito.ArgumentCaptor;
import org.mockito.Mock; import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension; import org.mockito.junit.jupiter.MockitoExtension;
@ -38,11 +40,10 @@ import org.thingsboard.server.common.data.msg.TbNodeConnectionType;
import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgMetaData; import org.thingsboard.server.common.msg.TbMsgMetaData;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List; import java.util.List;
import java.util.Map; import java.util.Map;
import java.util.UUID; import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
import java.util.stream.Stream; import java.util.stream.Stream;
@ -52,16 +53,18 @@ import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.eq; import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.BDDMockito.given; import static org.mockito.BDDMockito.given;
import static org.mockito.BDDMockito.never;
import static org.mockito.BDDMockito.spy; import static org.mockito.BDDMockito.spy;
import static org.mockito.BDDMockito.then; import static org.mockito.BDDMockito.then;
import static org.mockito.BDDMockito.times;
import static org.mockito.BDDMockito.willAnswer; import static org.mockito.BDDMockito.willAnswer;
@ExtendWith(MockitoExtension.class) @ExtendWith(MockitoExtension.class)
public class TbMsgDelayNodeTest extends AbstractRuleNodeUpgradeTest { public class TbMsgDelayNodeTest extends AbstractRuleNodeUpgradeTest {
private final DeviceId DEVICE_ID = new DeviceId(UUID.fromString("20107cf0-1c5e-4ac4-8131-7c466c955a7c")); private final DeviceId DEVICE_ID = new DeviceId(UUID.fromString("20107cf0-1c5e-4ac4-8131-7c466c955a7c"));
private final RuleNodeId RULE_NODE_ID = new RuleNodeId(UUID.fromString("1be24225-b669-4b26-ab7e-083aaa82d0a0"));
private final List<TimeUnit> supportedTimeUnits = List.of(TimeUnit.SECONDS, TimeUnit.MINUTES, TimeUnit.HOURS);
private final String supportedTimeUnitsStr = String.join(",", TimeUnit.SECONDS.name(), TimeUnit.MINUTES.name(), TimeUnit.HOURS.name());
private TbMsgDelayNode node; private TbMsgDelayNode node;
private TbMsgDelayNodeConfiguration config; private TbMsgDelayNodeConfiguration config;
@ -87,6 +90,40 @@ public class TbMsgDelayNodeTest extends AbstractRuleNodeUpgradeTest {
assertThatNoException().isThrownBy(() -> node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)))); assertThatNoException().isThrownBy(() -> node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))));
} }
@ParameterizedTest
@ValueSource(ints = {-1, 0, 5000000})
public void givenInvalidMaxPendingMsgsValue_whenInit_thenThrowsException() {
config.setMaxPendingMsgs(-1);
assertThatThrownBy(() -> node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))))
.isInstanceOf(TbNodeException.class)
.hasMessage("Maximum pending messages should be in a range from 1 to 100000.")
.extracting(e -> ((TbNodeException) e).isUnrecoverable())
.isEqualTo(true);
}
@Test
public void givenPeriodIsNull_whenInit_thenThrowsException() {
config.setPeriod(null);
assertThatThrownBy(() -> node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))))
.isInstanceOf(TbNodeException.class)
.hasMessage("Period should be specified.")
.extracting(e -> ((TbNodeException) e).isUnrecoverable())
.isEqualTo(true);
}
@Test
public void givenTimeUnitIsNull_whenInit_thenThrowsException() {
config.setTimeUnit(null);
assertThatThrownBy(() -> node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))))
.isInstanceOf(TbNodeException.class)
.hasMessage("Time unit should be specified.")
.extracting(e -> ((TbNodeException) e).isUnrecoverable())
.isEqualTo(true);
}
@ParameterizedTest @ParameterizedTest
@MethodSource @MethodSource
public void givenPeriodValueAndPeriodTimeUnitPatterns_whenOnMsg_thenTellSelfTickMsgAndEnqueueForTellNext( public void givenPeriodValueAndPeriodTimeUnitPatterns_whenOnMsg_thenTellSelfTickMsgAndEnqueueForTellNext(
@ -96,41 +133,23 @@ public class TbMsgDelayNodeTest extends AbstractRuleNodeUpgradeTest {
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))); node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
var ruleNodeId = new RuleNodeId(UUID.fromString("e8172ef8-bf91-4821-b9f5-ccd7b865e418")); var msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, metaData, data);
given(ctxMock.getSelfId()).willReturn(ruleNodeId); var tickMsg = TbMsg.newMsg(TbMsgType.DELAY_TIMEOUT_SELF_MSG, RULE_NODE_ID, TbMsgMetaData.EMPTY, msg.getId().toString());
given(ctxMock.newMsg(any(), any(TbMsgType.class), any(), any(), any(), any())).willReturn(tickMsg);
given(ctxMock.getSelfId()).willReturn(RULE_NODE_ID);
willAnswer(invocation -> { willAnswer(invocation -> {
node.onMsg(ctxMock, invocation.getArgument(0)); node.onMsg(ctxMock, invocation.getArgument(0));
return null; return null;
}).given(ctxMock).tellSelf(any(TbMsg.class), any(Long.class)); }).given(ctxMock).tellSelf(any(TbMsg.class), any(Long.class));
List<TbMsg> incomingMsgs = new ArrayList<>(); node.onMsg(ctxMock, msg);
List<TbMsg> tickMsgs = new ArrayList<>();
for (int i = 0; i < 9; i++) {
var msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, metaData, data);
incomingMsgs.add(msg);
var tickMsg = TbMsg.newMsg(TbMsgType.DELAY_TIMEOUT_SELF_MSG, ruleNodeId, TbMsgMetaData.EMPTY, msg.getId().toString());
tickMsgs.add(tickMsg);
given(ctxMock.newMsg(any(), any(TbMsgType.class), any(), any(), any(), eq(msg.getId().toString()))).willReturn(tickMsg);
}
incomingMsgs.forEach(msg -> node.onMsg(ctxMock, msg)); then(ctxMock).should().tellSelf(tickMsg, expectedDelay);
then(ctxMock).should().ack(msg);
incomingMsgs.forEach(incomingMsg -> { ArgumentCaptor<TbMsg> actualMsg = ArgumentCaptor.forClass(TbMsg.class);
then(ctxMock).should().newMsg(incomingMsg.getQueueName(), TbMsgType.DELAY_TIMEOUT_SELF_MSG, ruleNodeId, incomingMsg.getCustomerId(), TbMsgMetaData.EMPTY, incomingMsg.getId().toString()); then(ctxMock).should().enqueueForTellNext(actualMsg.capture(), eq(TbNodeConnectionType.SUCCESS));
then(ctxMock).should().ack(incomingMsg); assertThat(actualMsg.getValue()).usingRecursiveComparison().ignoringFields("id", "ts").isEqualTo(msg);
});
tickMsgs.forEach(tickMsg -> {
then(ctxMock).should().tellSelf(tickMsg, expectedDelay);
then(node).should().onMsg(ctxMock, tickMsg);
});
var actualMsgsCaptor = ArgumentCaptor.forClass(TbMsg.class);
then(ctxMock).should(times(9)).enqueueForTellNext(actualMsgsCaptor.capture(), eq(TbNodeConnectionType.SUCCESS));
var actualMsgs = actualMsgsCaptor.getAllValues();
for (int i = 0; i < 9; i++) {
var actualMsg = actualMsgs.get(i);
then(ctxMock).should().enqueueForTellNext(actualMsg, TbNodeConnectionType.SUCCESS);
assertThat(actualMsg).usingRecursiveComparison().ignoringFields("id", "ts").isEqualTo(incomingMsgs.get(i));
}
} }
private static Stream<Arguments> givenPeriodValueAndPeriodTimeUnitPatterns_whenOnMsg_thenTellSelfTickMsgAndEnqueueForTellNext() { private static Stream<Arguments> givenPeriodValueAndPeriodTimeUnitPatterns_whenOnMsg_thenTellSelfTickMsgAndEnqueueForTellNext() {
@ -138,28 +157,38 @@ public class TbMsgDelayNodeTest extends AbstractRuleNodeUpgradeTest {
Arguments.of("1", "HOURS", TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT, TimeUnit.HOURS.toMillis(1L)), Arguments.of("1", "HOURS", TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT, TimeUnit.HOURS.toMillis(1L)),
Arguments.of("${md-period}", "${md-time-unit}", Arguments.of("${md-period}", "${md-time-unit}",
new TbMsgMetaData(Map.of( new TbMsgMetaData(Map.of(
"md-period", "5", "md-period", "5",
"md-time-unit", "MINUTES" "md-time-unit", "MINUTES"
)), TbMsg.EMPTY_JSON_OBJECT, TimeUnit.MINUTES.toMillis(5L)), )), TbMsg.EMPTY_JSON_OBJECT, TimeUnit.MINUTES.toMillis(5L)),
Arguments.of("$[msg-period]", "$[msg-time-unit]", TbMsgMetaData.EMPTY, Arguments.of("$[msg-period]", "$[msg-time-unit]", TbMsgMetaData.EMPTY,
"{\"msg-period\":10,\"msg-time-unit\":\"SECONDS\"}", TimeUnit.SECONDS.toMillis(10L)) "{\"msg-period\":10,\"msg-time-unit\":\"SECONDS\"}", TimeUnit.SECONDS.toMillis(10L))
); );
} }
@ParameterizedTest
@EnumSource(TimeUnit.class)
public void givenTimeUnit_whenOnMsg_thenVerify(TimeUnit timeUnit) throws TbNodeException {
config.setTimeUnit(timeUnit.name());
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
var msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT);
if (supportedTimeUnits.contains(timeUnit)) {
assertThatNoException().isThrownBy(() -> node.onMsg(ctxMock, msg));
} else {
assertThatThrownBy(() -> node.onMsg(ctxMock, msg))
.isInstanceOf(RuntimeException.class)
.hasMessage("Time unit '" + timeUnit + "' is not supported! Only " + supportedTimeUnitsStr + " are supported.");
}
}
@Test @Test
public void givenPeriodIsUnparsable_whenOnMsg_thenThrowsException() throws TbNodeException { public void givenPeriodIsUnparsable_whenOnMsg_thenThrowsException() throws TbNodeException {
config.setPeriod("five"); config.setPeriod("five");
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))); node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
var msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT); var msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT);
var ruleNodeId = new RuleNodeId(UUID.fromString("5236e9b9-1e29-4b95-b219-7043ff8f0414"));
var tickMsg = TbMsg.newMsg(TbMsgType.DELAY_TIMEOUT_SELF_MSG, ruleNodeId, TbMsgMetaData.EMPTY, msg.getId().toString());
given(ctxMock.getSelfId()).willReturn(ruleNodeId);
given(ctxMock.newMsg(any(), any(TbMsgType.class), any(), any(), any(), any())).willReturn(tickMsg);
assertThatThrownBy(() -> node.onMsg(ctxMock, msg)) assertThatThrownBy(() -> node.onMsg(ctxMock, msg))
.isInstanceOf(RuntimeException.class) .isInstanceOf(NumberFormatException.class)
.hasMessage("Can't parse period value : five"); .hasMessage("Can't parse period value : five");
} }
@ -169,58 +198,30 @@ public class TbMsgDelayNodeTest extends AbstractRuleNodeUpgradeTest {
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))); node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
var msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT); var msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT);
var ruleNodeId = new RuleNodeId(UUID.fromString("0210d69a-247f-4488-a242-dd3244b55088"));
var tickMsg = TbMsg.newMsg(TbMsgType.DELAY_TIMEOUT_SELF_MSG, ruleNodeId, TbMsgMetaData.EMPTY, msg.getId().toString());
given(ctxMock.getSelfId()).willReturn(ruleNodeId);
given(ctxMock.newMsg(any(), any(TbMsgType.class), any(), any(), any(), any())).willReturn(tickMsg);
assertThatThrownBy(() -> node.onMsg(ctxMock, msg)) assertThatThrownBy(() -> node.onMsg(ctxMock, msg))
.isInstanceOf(RuntimeException.class) .isInstanceOf(IllegalArgumentException.class)
.hasMessage("Invalid value for period time unit : sec"); .hasMessage("Invalid value for period time unit : sec");
} }
@Test @Test
public void givenMaxLimitOfPendingMsgsReached_whenOnMsg_thenTellFailure() throws TbNodeException { public void givenMaxLimitOfPendingMsgsReached_whenOnMsg_thenTellFailure() throws TbNodeException {
int maxPendingMsgs = 5; config.setMaxPendingMsgs(1);
config.setMaxPendingMsgs(maxPendingMsgs);
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))); node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
var msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT);
RuleNodeId ruleNodeId = new RuleNodeId(UUID.fromString("d1440f09-ca81-41f3-b67e-1495aee87dc6")); for (int i = 0; i < 2; i++) {
given(ctxMock.getSelfId()).willReturn(ruleNodeId); node.onMsg(ctxMock, msg);
List<TbMsg> incomingMsgs = new ArrayList<>();
List<TbMsg> tickMsgs = new ArrayList<>();
for (int i = 0; i < 6; i++) {
var msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT);
incomingMsgs.add(msg);
var tickMsg = TbMsg.newMsg(TbMsgType.DELAY_TIMEOUT_SELF_MSG, ruleNodeId, TbMsgMetaData.EMPTY, msg.getId().toString());
tickMsgs.add(tickMsg);
}
for (int i = 0; i < maxPendingMsgs; i++) {
given(ctxMock.newMsg(any(), any(TbMsgType.class), any(), any(), any(), eq(incomingMsgs.get(i).getId().toString()))).willReturn(tickMsgs.get(i));
} }
incomingMsgs.forEach(msg -> node.onMsg(ctxMock, msg));
var lastMsg = incomingMsgs.remove(maxPendingMsgs);
incomingMsgs.forEach(incomingMsg -> {
then(ctxMock).should().newMsg(incomingMsg.getQueueName(), TbMsgType.DELAY_TIMEOUT_SELF_MSG, ruleNodeId, incomingMsg.getCustomerId(), TbMsgMetaData.EMPTY, incomingMsg.getId().toString());
then(ctxMock).should().ack(incomingMsg);
});
var tickForLastMsg = tickMsgs.remove(maxPendingMsgs);
tickMsgs.forEach(tickMsg -> then(ctxMock).should().tellSelf(tickMsg, TimeUnit.SECONDS.toMillis(60L)));
then(ctxMock).should(never()).tellSelf(eq(tickForLastMsg), any(Long.class));
ArgumentCaptor<Throwable> throwable = ArgumentCaptor.forClass(Throwable.class); ArgumentCaptor<Throwable> throwable = ArgumentCaptor.forClass(Throwable.class);
then(ctxMock).should().tellFailure(eq(lastMsg), throwable.capture()); then(ctxMock).should().tellFailure(eq(msg), throwable.capture());
assertThat(throwable.getValue()).isInstanceOf(RuntimeException.class).hasMessage("Max limit of pending messages reached!"); assertThat(throwable.getValue()).isInstanceOf(RuntimeException.class).hasMessage("Max limit of pending messages reached!");
} }
@Test @Test
public void verifyDestroyMethod() { public void verifyDestroyMethod() {
var msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT); var msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT);
var pendingMsgs = new HashMap<UUID, TbMsg>(); var pendingMsgs = new ConcurrentHashMap<>();
pendingMsgs.put(UUID.fromString("321f0301-9bed-4e7d-b92f-a978f53ec5d6"), msg); pendingMsgs.put(UUID.fromString("321f0301-9bed-4e7d-b92f-a978f53ec5d6"), msg);
ReflectionTestUtils.setField(node, "pendingMsgs", pendingMsgs); ReflectionTestUtils.setField(node, "pendingMsgs", pendingMsgs);
var actualPendingMsgs = (Map<UUID, TbMsg>) ReflectionTestUtils.getField(node, "pendingMsgs"); var actualPendingMsgs = (Map<UUID, TbMsg>) ReflectionTestUtils.getField(node, "pendingMsgs");
@ -233,61 +234,79 @@ public class TbMsgDelayNodeTest extends AbstractRuleNodeUpgradeTest {
private static Stream<Arguments> givenFromVersionAndConfig_whenUpgrade_thenVerifyHasChangesAndConfig() { private static Stream<Arguments> givenFromVersionAndConfig_whenUpgrade_thenVerifyHasChangesAndConfig() {
return Stream.of( return Stream.of(
// config for version 1 with upgrade from version 0 // config for version 1 with upgrade from version 0 (useMetadataPeriodInSecondsPatterns does not exist)
Arguments.of(0, Arguments.of(0,
""" """
{ {
"periodInSeconds": 60, "periodInSeconds": 13,
"maxPendingMsgs": 1000, "maxPendingMsgs": 1000,
"periodInSecondsPattern": null, "periodInSecondsPattern": "17"
"useMetadataPeriodInSecondsPatterns": false }
} """,
""",
true, true,
""" """
{ {
"period": 60, "period": "60",
"timeUnit": "SECONDS", "timeUnit": "SECONDS",
"maxPendingMsgs": 1000 "maxPendingMsgs": 1000
} }
"""
),
// config for version 1 with upgrade from version 0 (useMetadataPeriodInSecondsPatterns is false)
Arguments.of(0,
""" """
{
"periodInSeconds": 60,
"maxPendingMsgs": 1000,
"periodInSecondsPattern": null,
"useMetadataPeriodInSecondsPatterns": false
}
""",
true,
"""
{
"period": "60",
"timeUnit": "SECONDS",
"maxPendingMsgs": 1000
}
"""
), ),
// config for version 1 with upgrade from version 0 (useMetadataPeriodInSecondsPattern is true) // config for version 1 with upgrade from version 0 (useMetadataPeriodInSecondsPattern is true)
Arguments.of(0, Arguments.of(0,
""" """
{ {
"periodInSeconds": 60, "periodInSeconds": 60,
"maxPendingMsgs": 1000, "maxPendingMsgs": 1000,
"periodInSecondsPattern": "${period-pattern}", "periodInSecondsPattern": "${period-pattern}",
"useMetadataPeriodInSecondsPatterns": true "useMetadataPeriodInSecondsPatterns": true
} }
""", """,
true, true,
""" """
{ {
"period": "${period-pattern}", "period": "${period-pattern}",
"timeUnit": "SECONDS", "timeUnit": "SECONDS",
"maxPendingMsgs": 1000 "maxPendingMsgs": 1000
} }
""" """
), ),
// config for version 1 with upgrade from version 0 (hasChanges is false) // config for version 1 with upgrade from version 0 (hasChanges is false)
Arguments.of(0, Arguments.of(0,
""" """
{ {
"period": "${period-pattern}", "period": "${period-pattern}",
"timeUnit": "SECONDS", "timeUnit": "SECONDS",
"maxPendingMsgs": 1000 "maxPendingMsgs": 1000
} }
""", """,
false, false,
""" """
{ {
"period": "${period-pattern}", "period": "${period-pattern}",
"timeUnit": "SECONDS", "timeUnit": "SECONDS",
"maxPendingMsgs": 1000 "maxPendingMsgs": 1000
} }
""" """
) )
); );
} }

Loading…
Cancel
Save