Browse Source

added checks on init() and provided upgrade for existing nodes

pull/11087/head
IrynaMatveieva 2 years ago
parent
commit
b149957e43
  1. 92
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNode.java
  2. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNodeConfiguration.java
  3. 211
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNodeTest.java

92
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNode.java

@ -31,7 +31,6 @@ import org.thingsboard.rule.engine.api.TbNode;
import org.thingsboard.rule.engine.api.TbNodeConfiguration;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.rule.engine.api.util.TbNodeUtils;
import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.kv.Aggregation;
import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery;
import org.thingsboard.server.common.data.kv.ReadTsKvQuery;
@ -75,27 +74,43 @@ public class TbGetTelemetryNode implements TbNode {
this.config = TbNodeUtils.convert(configuration, TbGetTelemetryNodeConfiguration.class);
tsKeyNames = config.getLatestTsKeyNames();
if (tsKeyNames.isEmpty()) {
throw new TbNodeException("Telemetry is not selected!", true);
throw new TbNodeException("Telemetry should be specified!", true);
}
if (config.getFetchMode() == null) {
fetchMode = config.getFetchMode();
if (fetchMode == null) {
throw new TbNodeException("FetchMode should be specified!", true);
}
limit = FetchMode.ALL.equals(config.getFetchMode()) ? validateLimit(config.getLimit()) : 1;
fetchMode = config.getFetchMode();
orderBy = getOrderBy();
aggregation = parseAggregationConfig(config.getAggregation());
}
Aggregation parseAggregationConfig(String aggName) {
if (StringUtils.isEmpty(aggName) || !fetchMode.equals(FetchMode.ALL)) {
return Aggregation.NONE;
switch (fetchMode) {
case ALL:
limit = validateLimit(config.getLimit());
if (config.getOrderBy() == null) {
throw new TbNodeException("OrderBy should be specified!", true);
}
orderBy = config.getOrderBy();
if (config.getAggregation() == null) {
throw new TbNodeException("Aggregation should be specified!", true);
}
aggregation = config.getAggregation();
break;
case FIRST:
limit = 1;
orderBy = Direction.ASC;
aggregation = Aggregation.NONE;
break;
case LAST:
limit = 1;
orderBy = Direction.DESC;
aggregation = Aggregation.NONE;
break;
}
return Aggregation.valueOf(aggName);
}
@Override
public void onMsg(TbContext ctx, TbMsg msg) {
Interval interval = getInterval(msg);
if (interval.getStartTs() > interval.getEndTs()) {
throw new RuntimeException("Interval start should be less than Interval end");
}
List<String> keys = TbNodeUtils.processPatterns(tsKeyNames, msg);
ListenableFuture<List<TsKvEntry>> list = ctx.getTimeseriesService().findAll(ctx.getTenantId(), msg.getOriginator(), buildQueries(interval, keys));
DonAsynchron.withCallback(list, data -> {
@ -115,22 +130,6 @@ public class TbGetTelemetryNode implements TbNode {
.collect(Collectors.toList());
}
private Direction getOrderBy() throws TbNodeException {
switch (fetchMode) {
case ALL:
if (config.getOrderBy() == null) {
throw new TbNodeException("OrderBy should be specified!", true);
}
return config.getOrderBy();
case FIRST:
return Direction.ASC;
case LAST:
return Direction.DESC;
default:
throw new TbNodeException("FetchMode '" + fetchMode + "' is not supported.", true);
}
}
private TbMsgMetaData updateMetadata(List<TsKvEntry> entries, TbMsg msg, List<String> keys) {
ObjectNode resultNode = JacksonUtil.newObjectNode(JacksonUtil.ALLOW_UNQUOTED_FIELD_NAMES_MAPPER);
if (FetchMode.ALL.equals(fetchMode)) {
@ -234,9 +233,38 @@ public class TbGetTelemetryNode implements TbNode {
boolean hasChanges = false;
switch (fromVersion) {
case 0 -> {
if (oldConfiguration.has("orderBy") && oldConfiguration.get("orderBy").isNull()) {
((ObjectNode) oldConfiguration).put("orderBy", Direction.ASC.name());
hasChanges = true;
if (oldConfiguration.hasNonNull("fetchMode")) {
String fetchMode = oldConfiguration.get("fetchMode").asText();
switch (fetchMode) {
case "FIRST":
((ObjectNode) oldConfiguration).put("orderBy", Direction.ASC.name());
((ObjectNode) oldConfiguration).put("aggregation", Aggregation.NONE.name());
hasChanges = true;
break;
case "LAST":
((ObjectNode) oldConfiguration).put("orderBy", Direction.DESC.name());
((ObjectNode) oldConfiguration).put("aggregation", Aggregation.NONE.name());
hasChanges = true;
break;
case "ALL":
if (oldConfiguration.has("orderBy") &&
(oldConfiguration.get("orderBy").isNull() || oldConfiguration.get("orderBy").asText().isEmpty())) {
((ObjectNode) oldConfiguration).put("orderBy", Direction.ASC.name());
hasChanges = true;
}
if (oldConfiguration.has("aggregation") &&
(oldConfiguration.get("aggregation").isNull() || oldConfiguration.get("aggregation").asText().isEmpty())) {
((ObjectNode) oldConfiguration).put("aggregation", Aggregation.NONE.name());
hasChanges = true;
}
break;
default:
((ObjectNode) oldConfiguration).put("fetchMode", FetchMode.LAST.name());
((ObjectNode) oldConfiguration).put("orderBy", Direction.DESC.name());
((ObjectNode) oldConfiguration).put("aggregation", Aggregation.NONE.name());
hasChanges = true;
break;
}
}
}
}

4
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNodeConfiguration.java

@ -44,7 +44,7 @@ public class TbGetTelemetryNodeConfiguration implements NodeConfiguration<TbGetT
private String endIntervalTimeUnit;
private FetchMode fetchMode; //FIRST, LAST, ALL
private Direction orderBy; //ASC, DESC
private String aggregation; //MIN, MAX, AVG, SUM, COUNT, NONE;
private Aggregation aggregation; //MIN, MAX, AVG, SUM, COUNT, NONE;
private int limit;
private List<String> latestTsKeyNames;
@ -62,7 +62,7 @@ public class TbGetTelemetryNodeConfiguration implements NodeConfiguration<TbGetT
configuration.setStartIntervalPattern("");
configuration.setEndIntervalPattern("");
configuration.setOrderBy(Direction.ASC);
configuration.setAggregation(Aggregation.NONE.name());
configuration.setAggregation(Aggregation.NONE);
configuration.setLimit(MAX_FETCH_SIZE);
return configuration;
}

211
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNodeTest.java

@ -89,49 +89,6 @@ public class TbGetTelemetryNodeTest extends AbstractRuleNodeUpgradeTest {
config.setLatestTsKeyNames(List.of("temperature"));
}
@Test
public void givenAggregationAsString_whenParseAggregation_thenReturnEnum() throws TbNodeException {
config.setFetchMode(FetchMode.ALL);
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
//compatibility with old configs without "aggregation" parameter
assertThat(node.parseAggregationConfig(null)).isEqualTo(Aggregation.NONE);
assertThat(node.parseAggregationConfig("")).isEqualTo(Aggregation.NONE);
//common values
assertThat(node.parseAggregationConfig("MIN")).isEqualTo(Aggregation.MIN);
assertThat(node.parseAggregationConfig("MAX")).isEqualTo(Aggregation.MAX);
assertThat(node.parseAggregationConfig("AVG")).isEqualTo(Aggregation.AVG);
assertThat(node.parseAggregationConfig("SUM")).isEqualTo(Aggregation.SUM);
assertThat(node.parseAggregationConfig("COUNT")).isEqualTo(Aggregation.COUNT);
assertThat(node.parseAggregationConfig("NONE")).isEqualTo(Aggregation.NONE);
//all possible values in future
for (Aggregation aggEnum : Aggregation.values()) {
assertThat(node.parseAggregationConfig(aggEnum.name())).isEqualTo(aggEnum);
}
}
@Test
public void givenAggregationWhiteSpace_whenParseAggregation_thenException() throws TbNodeException {
// GIVEN
config.setFetchMode(FetchMode.ALL);
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
// WHEN-THEN
assertThatThrownBy(() -> node.parseAggregationConfig(" ")).isInstanceOf(IllegalArgumentException.class);
}
@Test
public void givenAggregationIncorrect_whenParseAggregation_thenException() throws TbNodeException {
// GIVEN
config.setFetchMode(FetchMode.ALL);
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
// WHEN-THEN
assertThatThrownBy(() -> node.parseAggregationConfig("TOP")).isInstanceOf(IllegalArgumentException.class);
}
@Test
public void verifyDefaultConfig() {
config = new TbGetTelemetryNodeConfiguration().defaultConfiguration();
@ -144,7 +101,7 @@ public class TbGetTelemetryNodeTest extends AbstractRuleNodeUpgradeTest {
assertThat(config.getEndIntervalTimeUnit()).isEqualTo(TimeUnit.MINUTES.name());
assertThat(config.getFetchMode()).isEqualTo(FetchMode.FIRST);
assertThat(config.getOrderBy()).isEqualTo(Direction.ASC);
assertThat(config.getAggregation()).isEqualTo(Aggregation.NONE.name());
assertThat(config.getAggregation()).isEqualTo(Aggregation.NONE);
assertThat(config.getLimit()).isEqualTo(1000);
assertThat(config.getLatestTsKeyNames()).isEmpty();
}
@ -157,7 +114,7 @@ public class TbGetTelemetryNodeTest extends AbstractRuleNodeUpgradeTest {
// WHEN-THEN
assertThatThrownBy(() -> node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))))
.isInstanceOf(TbNodeException.class)
.hasMessage("Telemetry is not selected!")
.hasMessage("Telemetry should be specified!")
.extracting(e -> ((TbNodeException) e).isUnrecoverable())
.isEqualTo(true);
}
@ -181,7 +138,7 @@ public class TbGetTelemetryNodeTest extends AbstractRuleNodeUpgradeTest {
config.setFetchMode(FetchMode.ALL);
config.setOrderBy(null);
// THEN
// WHEN-THEN
assertThatThrownBy(() -> node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))))
.isInstanceOf(TbNodeException.class)
.hasMessage("OrderBy should be specified!")
@ -196,7 +153,7 @@ public class TbGetTelemetryNodeTest extends AbstractRuleNodeUpgradeTest {
config.setFetchMode(FetchMode.ALL);
config.setLimit(limit);
// THEN
// WHEN-THEN
assertThatThrownBy(() -> node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))))
.isInstanceOf(TbNodeException.class)
.hasMessage("Limit should be in a range from 2 to 1000.")
@ -204,6 +161,33 @@ public class TbGetTelemetryNodeTest extends AbstractRuleNodeUpgradeTest {
.isEqualTo(true);
}
@Test
public void givenFetchModeIsAllAndAggregationIsNull_whenInit_thenThrowsException() {
// GIVEN
config.setFetchMode(FetchMode.ALL);
config.setAggregation(null);
// WHEN-THEN
assertThatThrownBy(() -> node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))))
.isInstanceOf(TbNodeException.class)
.hasMessage("Aggregation should be specified!")
.extracting(e -> ((TbNodeException) e).isUnrecoverable())
.isEqualTo(true);
}
@Test
public void givenIntervalStartIsGreaterThanIntervalEnd_whenOnMsg_thenThrowsException() throws TbNodeException {
// GIVEN
config.setStartInterval(1);
config.setEndInterval(2);
// WHEN-THEN
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT);
assertThatThrownBy(() -> node.onMsg(ctxMock, msg))
.isInstanceOf(RuntimeException.class)
.hasMessage("Interval start should be less than Interval end");
}
@Test
public void givenUseMetadataIntervalPatternsIsTrue_whenOnMsg_thenVerifyStartAndEndTsInQuery() throws TbNodeException {
// GIVEN
@ -279,7 +263,7 @@ public class TbGetTelemetryNodeTest extends AbstractRuleNodeUpgradeTest {
@ParameterizedTest
@MethodSource
public void givenAggregation_whenOnMsg_thenVerifyAggregationStepInQuery(String aggregation, Consumer<ReadTsKvQuery> verifyAggregationStepInQuery) throws TbNodeException {
public void givenAggregation_whenOnMsg_thenVerifyAggregationStepInQuery(Aggregation aggregation, Consumer<ReadTsKvQuery> verifyAggregationStepInQuery) throws TbNodeException {
// GIVEN
config.setStartInterval(5);
config.setEndInterval(1);
@ -304,8 +288,8 @@ public class TbGetTelemetryNodeTest extends AbstractRuleNodeUpgradeTest {
private static Stream<Arguments> givenAggregation_whenOnMsg_thenVerifyAggregationStepInQuery() {
return Stream.of(
Arguments.of("", (Consumer<ReadTsKvQuery>) query -> assertThat(query.getInterval()).isEqualTo(1)),
Arguments.of("MIN", (Consumer<ReadTsKvQuery>) query -> assertThat(query.getInterval()).isEqualTo(query.getEndTs() - query.getStartTs()))
Arguments.of(Aggregation.NONE, (Consumer<ReadTsKvQuery>) query -> assertThat(query.getInterval()).isEqualTo(1)),
Arguments.of(Aggregation.AVG, (Consumer<ReadTsKvQuery>) query -> assertThat(query.getInterval()).isEqualTo(query.getEndTs() - query.getStartTs()))
);
}
@ -479,15 +463,84 @@ public class TbGetTelemetryNodeTest extends AbstractRuleNodeUpgradeTest {
private static Stream<Arguments> givenFromVersionAndConfig_whenUpgrade_thenVerifyHasChangesAndConfig() {
return Stream.of(
// config for version 0 (orderBy id null)
// config for version 0 (fetchMode is 'FIRST' and orderBy is 'INVALID_ORDER_BY' and aggregation is 'SUM')
Arguments.of(0,
"""
{
"latestTsKeyNames": [],
"aggregation": "SUM",
"fetchMode": "FIRST",
"orderBy": "INVALID_ORDER_BY",
"limit": 1000,
"useMetadataIntervalPatterns": false,
"startIntervalPattern": "",
"endIntervalPattern": "",
"startInterval": 2,
"startIntervalTimeUnit": "MINUTES",
"endInterval": 1,
"endIntervalTimeUnit": "MINUTES"
}
""",
true,
"""
{
"latestTsKeyNames": [],
"aggregation": "NONE",
"fetchMode": "FIRST",
"orderBy": "ASC",
"limit": 1000,
"useMetadataIntervalPatterns": false,
"startIntervalPattern": "",
"endIntervalPattern": "",
"startInterval": 2,
"startIntervalTimeUnit": "MINUTES",
"endInterval": 1,
"endIntervalTimeUnit": "MINUTES"
}
"""),
// config for version 0 (fetchMode is 'LAST' and orderBy is 'ASC' and aggregation is 'AVG')
Arguments.of(0,
"""
{
"latestTsKeyNames": [
],
"latestTsKeyNames": [],
"aggregation": "AVG",
"fetchMode": "LAST",
"orderBy": "ASC",
"limit": 1000,
"useMetadataIntervalPatterns": false,
"startIntervalPattern": "",
"endIntervalPattern": "",
"startInterval": 2,
"startIntervalTimeUnit": "MINUTES",
"endInterval": 1,
"endIntervalTimeUnit": "MINUTES"
}
""",
true,
"""
{
"latestTsKeyNames": [],
"aggregation": "NONE",
"fetchMode": "LAST",
"orderBy": "DESC",
"limit": 1000,
"useMetadataIntervalPatterns": false,
"startIntervalPattern": "",
"endIntervalPattern": "",
"startInterval": 2,
"startIntervalTimeUnit": "MINUTES",
"endInterval": 1,
"endIntervalTimeUnit": "MINUTES"
}
"""),
// config for version 0 (fetchMode is 'ALL' and orderBy is empty and aggregation is null)
Arguments.of(0,
"""
{
"latestTsKeyNames": [],
"aggregation": null,
"fetchMode": "ALL",
"orderBy": null,
"orderBy": "",
"limit": 1000,
"useMetadataIntervalPatterns": false,
"startIntervalPattern": "",
@ -501,8 +554,7 @@ public class TbGetTelemetryNodeTest extends AbstractRuleNodeUpgradeTest {
true,
"""
{
"latestTsKeyNames": [
],
"latestTsKeyNames": [],
"aggregation": "NONE",
"fetchMode": "ALL",
"orderBy": "ASC",
@ -516,13 +568,12 @@ public class TbGetTelemetryNodeTest extends AbstractRuleNodeUpgradeTest {
"endIntervalTimeUnit": "MINUTES"
}
"""),
// config for version 0 (orderBy is specified)
// config for version 0 (fetchMode is 'ALL' and orderBy is 'DESC' and aggregation is 'SUM')
Arguments.of(0,
"""
{
"latestTsKeyNames": [
],
"aggregation": "NONE",
"latestTsKeyNames": [],
"aggregation": "SUM",
"fetchMode": "ALL",
"orderBy": "DESC",
"limit": 1000,
@ -538,9 +589,8 @@ public class TbGetTelemetryNodeTest extends AbstractRuleNodeUpgradeTest {
false,
"""
{
"latestTsKeyNames": [
],
"aggregation": "NONE",
"latestTsKeyNames": [],
"aggregation": "SUM",
"fetchMode": "ALL",
"orderBy": "DESC",
"limit": 1000,
@ -552,6 +602,41 @@ public class TbGetTelemetryNodeTest extends AbstractRuleNodeUpgradeTest {
"endInterval": 1,
"endIntervalTimeUnit": "MINUTES"
}
"""),
// config for version 0 (fetchMode is 'INVALID_MODE' and orderBy is 'INVALID_ORDER_BY' and aggregation is 'INVALID_AGGREGATION')
Arguments.of(0,
"""
{
"latestTsKeyNames": [],
"aggregation": "INVALID_AGGREGATION",
"fetchMode": "INVALID_MODE",
"orderBy": "INVALID_ORDER_BY",
"limit": 1000,
"useMetadataIntervalPatterns": false,
"startIntervalPattern": "",
"endIntervalPattern": "",
"startInterval": 2,
"startIntervalTimeUnit": "MINUTES",
"endInterval": 1,
"endIntervalTimeUnit": "MINUTES"
}
""",
true,
"""
{
"latestTsKeyNames": [],
"aggregation": "NONE",
"fetchMode": "LAST",
"orderBy": "DESC",
"limit": 1000,
"useMetadataIntervalPatterns": false,
"startIntervalPattern": "",
"endIntervalPattern": "",
"startInterval": 2,
"startIntervalTimeUnit": "MINUTES",
"endInterval": 1,
"endIntervalTimeUnit": "MINUTES"
}
""")
);
}

Loading…
Cancel
Save