From b149957e439ec4cf0c1d02a96c2888e5233f15c9 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Thu, 15 Aug 2024 16:23:26 +0300 Subject: [PATCH] added checks on init() and provided upgrade for existing nodes --- .../engine/metadata/TbGetTelemetryNode.java | 92 +++++--- .../TbGetTelemetryNodeConfiguration.java | 4 +- .../metadata/TbGetTelemetryNodeTest.java | 211 ++++++++++++------ 3 files changed, 210 insertions(+), 97 deletions(-) diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNode.java index 5faf0e40e2..2b5190563d 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNode.java +++ b/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 keys = TbNodeUtils.processPatterns(tsKeyNames, msg); ListenableFuture> 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 entries, TbMsg msg, List 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; + } } } } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNodeConfiguration.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNodeConfiguration.java index 3bb9ff1013..704abeddf6 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNodeConfiguration.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNodeConfiguration.java @@ -44,7 +44,7 @@ public class TbGetTelemetryNodeConfiguration implements NodeConfiguration latestTsKeyNames; @@ -62,7 +62,7 @@ public class TbGetTelemetryNodeConfiguration implements NodeConfiguration 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 verifyAggregationStepInQuery) throws TbNodeException { + public void givenAggregation_whenOnMsg_thenVerifyAggregationStepInQuery(Aggregation aggregation, Consumer verifyAggregationStepInQuery) throws TbNodeException { // GIVEN config.setStartInterval(5); config.setEndInterval(1); @@ -304,8 +288,8 @@ public class TbGetTelemetryNodeTest extends AbstractRuleNodeUpgradeTest { private static Stream givenAggregation_whenOnMsg_thenVerifyAggregationStepInQuery() { return Stream.of( - Arguments.of("", (Consumer) query -> assertThat(query.getInterval()).isEqualTo(1)), - Arguments.of("MIN", (Consumer) query -> assertThat(query.getInterval()).isEqualTo(query.getEndTs() - query.getStartTs())) + Arguments.of(Aggregation.NONE, (Consumer) query -> assertThat(query.getInterval()).isEqualTo(1)), + Arguments.of(Aggregation.AVG, (Consumer) query -> assertThat(query.getInterval()).isEqualTo(query.getEndTs() - query.getStartTs())) ); } @@ -479,15 +463,84 @@ public class TbGetTelemetryNodeTest extends AbstractRuleNodeUpgradeTest { private static Stream 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" + } """) ); }