diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/CalculateDeltaNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/CalculateDeltaNode.java index 609888db84..7141d66f42 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/CalculateDeltaNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/CalculateDeltaNode.java @@ -32,6 +32,7 @@ import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.data.msg.TbMsgType; import org.thingsboard.server.common.data.msg.TbNodeConnectionType; import org.thingsboard.server.common.data.plugin.ComponentType; +import org.thingsboard.server.common.data.util.TbPair; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.dao.timeseries.TimeseriesService; @@ -46,7 +47,9 @@ import static org.thingsboard.common.util.DonAsynchron.withCallback; @Slf4j @RuleNode(type = ComponentType.ENRICHMENT, - name = "calculate delta", relationTypes = {TbNodeConnectionType.SUCCESS, TbNodeConnectionType.FAILURE, TbNodeConnectionType.OTHER}, + name = "calculate delta", + version = 1, + relationTypes = {TbNodeConnectionType.SUCCESS, TbNodeConnectionType.FAILURE, TbNodeConnectionType.OTHER}, configClazz = CalculateDeltaNodeConfiguration.class, nodeDescription = "Calculates delta and amount of time passed between previous timeseries key reading " + "and current value for this key from the incoming message", @@ -101,6 +104,11 @@ public class CalculateDeltaNode implements TbNode { return; } + if (config.isExcludeZeroDeltas() && delta.doubleValue() == 0) { + ctx.tellSuccess(msg); + return; + } + if (config.getRound() != null) { delta = delta.setScale(config.getRound(), RoundingMode.HALF_UP); } @@ -128,6 +136,23 @@ public class CalculateDeltaNode implements TbNode { } } + @Override + public TbPair upgrade(int fromVersion, JsonNode oldConfiguration) throws TbNodeException { + boolean hasChanges = false; + switch (fromVersion) { + case 0: + String excludeZeroDeltas = "excludeZeroDeltas"; + if (!oldConfiguration.has(excludeZeroDeltas)) { + hasChanges = true; + ((ObjectNode) oldConfiguration).put(excludeZeroDeltas, false); + } + break; + default: + break; + } + return new TbPair<>(hasChanges, oldConfiguration); + } + private ListenableFuture fetchLatestValueAsync(EntityId entityId) { return Futures.transform(timeseriesService.findLatest(ctx.getTenantId(), entityId, Collections.singletonList(config.getInputValueKey())), list -> extractValue(list.get(0)) diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/CalculateDeltaNodeConfiguration.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/CalculateDeltaNodeConfiguration.java index 0c4e6de556..0ae558718b 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/CalculateDeltaNodeConfiguration.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/CalculateDeltaNodeConfiguration.java @@ -30,6 +30,7 @@ public class CalculateDeltaNodeConfiguration implements NodeConfiguration CalculateDeltaTestConfig() { + return Stream.of( + // delta = 0, tell failure if delta is negative is set to true and exclude zero deltas is set to true so delta should filter out the message. + new CalculateDeltaTestConfig(true, true, 40, 40, (ctx, msg) -> { + verify(ctx).tellSuccess(eq(msg)); + verify(ctx).getDbCallbackExecutor(); + verifyNoMoreInteractions(ctx); + }), + // delta < 0, tell failure if delta is negative is set to true so it should throw exception. + new CalculateDeltaTestConfig(true, true, 41, 40, (ctx, msg) -> { + var errorCaptor = ArgumentCaptor.forClass(Throwable.class); + verify(ctx).tellFailure(eq(msg), errorCaptor.capture()); + verify(ctx).getDbCallbackExecutor(); + verifyNoMoreInteractions(ctx); + assertThat(errorCaptor.getValue()).isInstanceOf(IllegalArgumentException.class).hasMessage("Delta value is negative!"); + }), + // delta < 0, exclude zero deltas is set to true so it should return message with delta if delta is negative is set to false. + new CalculateDeltaTestConfig(false, true, 41, 40, (ctx, msg) -> { + var actualMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); + verify(ctx).tellSuccess(actualMsgCaptor.capture()); + verify(ctx).getDbCallbackExecutor(); + verifyNoMoreInteractions(ctx); + String expectedMsgData = "{\"temperature\":40.0,\"airPressure\":123,\"delta\":-1}"; + assertEquals(expectedMsgData, actualMsgCaptor.getValue().getData()); + }), + // delta = 0, tell failure if delta is negative is set to false and exclude zero deltas is set to true so delta should filter out the message. + new CalculateDeltaTestConfig(false, true, 40, 40, (ctx, msg) -> { + verify(ctx).tellSuccess(eq(msg)); + verify(ctx).getDbCallbackExecutor(); + verifyNoMoreInteractions(ctx); + }), + // delta > 0, exclude zero deltas is set to true so it should return message with delta. + new CalculateDeltaTestConfig(false, true, 39, 40, (ctx, msg) -> { + var actualMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); + verify(ctx).tellSuccess(actualMsgCaptor.capture()); + verify(ctx).getDbCallbackExecutor(); + verifyNoMoreInteractions(ctx); + String expectedMsgData = "{\"temperature\":40.0,\"airPressure\":123,\"delta\":1}"; + assertEquals(expectedMsgData, actualMsgCaptor.getValue().getData()); + }), + // delta > 0, exclude zero deltas is set to false so it should return message with delta. + new CalculateDeltaTestConfig(false, false, 39, 40, (ctx, msg) -> { + var actualMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); + verify(ctx).tellSuccess(actualMsgCaptor.capture()); + verify(ctx).getDbCallbackExecutor(); + verifyNoMoreInteractions(ctx); + String expectedMsgData = "{\"temperature\":40.0,\"airPressure\":123,\"delta\":1}"; + assertEquals(expectedMsgData, actualMsgCaptor.getValue().getData()); + }) + ); + } + + @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 verificationMethod; + } + private void mockFindLatest(TsKvEntry tsKvEntry) { when(ctxMock.getTenantId()).thenReturn(TENANT_ID); when(timeseriesServiceMock.findLatestSync( @@ -457,4 +553,24 @@ public class CalculateDeltaNodeTest { } + private static Stream givenFromVersionAndConfig_whenUpgrade_thenVerifyHasChangesAndConfig() { + return Stream.of( + // default config for version 0 + Arguments.of(0, + "{\"inputValueKey\":\"pulseCounter\",\"outputValueKey\":\"delta\",\"useCache\":true,\"addPeriodBetweenMsgs\":false, \"periodValueKey\":\"periodInMs\", \"round\":null,\"tellFailureIfDeltaIsNegative\":true}", + true, + "{\"inputValueKey\":\"pulseCounter\",\"outputValueKey\":\"delta\",\"useCache\":true,\"addPeriodBetweenMsgs\":false, \"periodValueKey\":\"periodInMs\", \"round\":null,\"tellFailureIfDeltaIsNegative\":true, \"excludeZeroDeltas\":false}"), + // default config for version 1 with upgrade from version 0 + Arguments.of(1, + "{\"inputValueKey\":\"pulseCounter\",\"outputValueKey\":\"delta\",\"useCache\":true,\"addPeriodBetweenMsgs\":false, \"periodValueKey\":\"periodInMs\", \"round\":null,\"tellFailureIfDeltaIsNegative\":true, \"excludeZeroDeltas\":false}", + false, + "{\"inputValueKey\":\"pulseCounter\",\"outputValueKey\":\"delta\",\"useCache\":true,\"addPeriodBetweenMsgs\":false, \"periodValueKey\":\"periodInMs\", \"round\":null,\"tellFailureIfDeltaIsNegative\":true, \"excludeZeroDeltas\":false}") + ); + + } + + @Override + protected TbNode getTestNode() { + return node; + } }