From 6afb3f25aaaa0acff4f63a12934606ac77df4b5c Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Tue, 2 Jul 2024 17:30:17 +0300 Subject: [PATCH 01/12] added tests for the delay node --- .../rule/engine/delay/TbMsgDelayNodeTest.java | 213 ++++++++++++++++++ 1 file changed, 213 insertions(+) create mode 100644 rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/delay/TbMsgDelayNodeTest.java diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/delay/TbMsgDelayNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/delay/TbMsgDelayNodeTest.java new file mode 100644 index 0000000000..18d3ffc2c5 --- /dev/null +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/delay/TbMsgDelayNodeTest.java @@ -0,0 +1,213 @@ +/** + * Copyright © 2016-2024 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.rule.engine.delay; + +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.mockito.ArgumentCaptor; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; +import org.springframework.test.util.ReflectionTestUtils; +import org.thingsboard.common.util.JacksonUtil; +import org.thingsboard.rule.engine.api.TbContext; +import org.thingsboard.rule.engine.api.TbNodeConfiguration; +import org.thingsboard.rule.engine.api.TbNodeException; +import org.thingsboard.server.common.data.id.DeviceId; +import org.thingsboard.server.common.data.id.RuleNodeId; +import org.thingsboard.server.common.data.msg.TbMsgType; +import org.thingsboard.server.common.data.msg.TbNodeConnectionType; +import org.thingsboard.server.common.msg.TbMsg; +import org.thingsboard.server.common.msg.TbMsgMetaData; + +import java.util.HashMap; +import java.util.Map; +import java.util.UUID; +import java.util.stream.Stream; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatNoException; +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.ArgumentMatchers.isNull; +import static org.mockito.BDDMockito.given; +import static org.mockito.BDDMockito.spy; +import static org.mockito.BDDMockito.then; +import static org.mockito.BDDMockito.times; +import static org.mockito.BDDMockito.willAnswer; + +@ExtendWith(MockitoExtension.class) +public class TbMsgDelayNodeTest { + + private final DeviceId DEVICE_ID = new DeviceId(UUID.fromString("20107cf0-1c5e-4ac4-8131-7c466c955a7c")); + + private TbMsgDelayNode node; + private TbMsgDelayNodeConfiguration config; + + @Mock + private TbContext ctxMock; + + @BeforeEach + public void setUp() { + node = spy(new TbMsgDelayNode()); + config = new TbMsgDelayNodeConfiguration().defaultConfiguration(); + } + + @Test + public void verifyDefaultConfig() { + assertThat(config.getPeriodInSeconds()).isEqualTo(60); + assertThat(config.getMaxPendingMsgs()).isEqualTo(1000); + assertThat(config.isUseMetadataPeriodInSecondsPatterns()).isFalse(); + assertThat(config.getPeriodInSecondsPattern()).isNull(); + } + + @Test + public void givenDefaultConfig_whenInit_thenOk() { + assertThatNoException().isThrownBy(() -> node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)))); + } + + @ParameterizedTest + @MethodSource + public void givenPeriodInSecondsPattern_whenOnMsg_thenTellSelfTickMsgAndEnqueueForTellNext( + String periodInSecondsPattern, TbMsgMetaData metaData, String data, long expectedDelay) throws TbNodeException { + config.setUseMetadataPeriodInSecondsPatterns(true); + config.setPeriodInSecondsPattern(periodInSecondsPattern); + node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))); + + TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, metaData, data); + RuleNodeId ruleNodeId = new RuleNodeId(UUID.fromString("5236e9b9-1e29-4b95-b219-7043ff8f0414")); + TbMsg 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); + willAnswer(invocation -> { + node.onMsg(ctxMock, invocation.getArgument(0)); + return null; + }).given(ctxMock).tellSelf(tickMsg, expectedDelay); + + node.onMsg(ctxMock, msg); + + then(ctxMock).should().newMsg(null, TbMsgType.DELAY_TIMEOUT_SELF_MSG, ruleNodeId, null, TbMsgMetaData.EMPTY, msg.getId().toString()); + then(ctxMock).should().tellSelf(tickMsg, expectedDelay); + then(ctxMock).should().ack(msg); + then(node).should().onMsg(ctxMock, tickMsg); + ArgumentCaptor actualMsg = ArgumentCaptor.forClass(TbMsg.class); + then(ctxMock).should().enqueueForTellNext(actualMsg.capture(), eq(TbNodeConnectionType.SUCCESS)); + assertThat(actualMsg.getValue()).usingRecursiveComparison().ignoringFields("id", "ts").isEqualTo(msg); + } + + private static Stream givenPeriodInSecondsPattern_whenOnMsg_thenTellSelfTickMsgAndEnqueueForTellNext() { + return Stream.of( + Arguments.of("${md-period-in-seconds}", new TbMsgMetaData(Map.of("md-period-in-seconds", "10")), TbMsg.EMPTY_JSON_OBJECT, 10000L), + Arguments.of("$[msg-period-in-seconds]", TbMsgMetaData.EMPTY, "{\"msg-period-in-seconds\":5}", 5000L) + ); + } + + @Test + public void givenPeriodInSecondsPatternIsUnparsable_whenOnMsg_thenThrowsException() throws TbNodeException { + config.setUseMetadataPeriodInSecondsPatterns(true); + config.setPeriodInSecondsPattern("$[msg-period-in-seconds]"); + node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))); + + TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, TbMsgMetaData.EMPTY, "{\"msg-period-in-seconds\":\"five\"}"); + RuleNodeId ruleNodeId = new RuleNodeId(UUID.fromString("5236e9b9-1e29-4b95-b219-7043ff8f0414")); + TbMsg 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)) + .isInstanceOf(RuntimeException.class) + .hasMessage("Can't parse period in seconds from metadata using pattern: $[msg-period-in-seconds]"); + } + + @Test + public void givenMaxLimitOfPendingMsgsReached_whenOnMsg_thenTellFailure() throws TbNodeException { + int maxPendingMsgs = 5; + config.setMaxPendingMsgs(maxPendingMsgs); + + node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))); + + RuleNodeId ruleNodeId = new RuleNodeId(UUID.fromString("d1440f09-ca81-41f3-b67e-1495aee87dc6")); + String msgId = "e38c87c5-8916-4bb0-b448-b8d08ad639df"; + TbMsg tickMsg = TbMsg.newMsg(TbMsgType.DELAY_TIMEOUT_SELF_MSG, ruleNodeId, TbMsgMetaData.EMPTY, msgId); + + given(ctxMock.getSelfId()).willReturn(ruleNodeId); + given(ctxMock.newMsg(any(), any(TbMsgType.class), any(), any(), any(), any())).willReturn(tickMsg); + + for (int i = 0; i < 6; i++) { + TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, TbMsgMetaData.EMPTY, "msg : " + (i + 1)); + node.onMsg(ctxMock, msg); + } + + then(ctxMock).should(times(maxPendingMsgs)).newMsg(isNull(), eq(TbMsgType.DELAY_TIMEOUT_SELF_MSG), eq(ruleNodeId), isNull(), eq(TbMsgMetaData.EMPTY), any()); + then(ctxMock).should(times(maxPendingMsgs)).tellSelf(tickMsg, 60000L); + then(ctxMock).should(times(maxPendingMsgs)).ack(any()); + ArgumentCaptor actualMsg = ArgumentCaptor.forClass(TbMsg.class); + ArgumentCaptor throwable = ArgumentCaptor.forClass(Throwable.class); + then(ctxMock).should().tellFailure(actualMsg.capture(), throwable.capture()); + TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, TbMsgMetaData.EMPTY, "msg : 6"); + assertThat(actualMsg.getValue()).usingRecursiveComparison().ignoringFields("id", "ts").isEqualTo(msg); + assertThat(throwable.getValue()).isInstanceOf(RuntimeException.class).hasMessage("Max limit of pending messages reached!"); + } + + @Test + public void givenNumberOfMsgsMoreThenMaxPendingMsgs_whenOnMsg_thenTellSelfTickMsgAndEnqueueForTellNext() throws TbNodeException { + config.setMaxPendingMsgs(3); + config.setPeriodInSeconds(1); + + node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))); + + RuleNodeId ruleNodeId = new RuleNodeId(UUID.fromString("e8172ef8-bf91-4821-b9f5-ccd7b865e418")); + TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT); + TbMsg 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); + willAnswer(invocation -> { + node.onMsg(ctxMock, invocation.getArgument(0)); + return null; + }).given(ctxMock).tellSelf(tickMsg, 1000L); + + for (int i = 0; i < 9; i++) { + node.onMsg(ctxMock, msg); + } + + then(ctxMock).should(times(9)).newMsg(null, TbMsgType.DELAY_TIMEOUT_SELF_MSG, ruleNodeId, null, TbMsgMetaData.EMPTY, msg.getId().toString()); + then(ctxMock).should(times(9)).tellSelf(tickMsg, 1000L); + then(ctxMock).should(times(9)).ack(msg); + then(node).should(times(9)).onMsg(ctxMock, tickMsg); + then(ctxMock).should(times(9)).enqueueForTellNext(any(TbMsg.class), eq(TbNodeConnectionType.SUCCESS)); + } + + @Test + public void verifyDestroyMethod() { + TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT); + var pendingMsgs = new HashMap(); + pendingMsgs.put(UUID.fromString("321f0301-9bed-4e7d-b92f-a978f53ec5d6"), msg); + ReflectionTestUtils.setField(node, "pendingMsgs", pendingMsgs); + var actualPendingMsgs = (Map) ReflectionTestUtils.getField(node, "pendingMsgs"); + assertThat(actualPendingMsgs).isEqualTo(pendingMsgs); + + node.destroy(); + + assertThat(actualPendingMsgs).isEmpty(); + } +} From 1e703f0264388c486a9c67d74680c6e94a73034e Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Wed, 3 Jul 2024 16:34:27 +0300 Subject: [PATCH 02/12] added ability to set time unit using pattern --- .../rule/engine/delay/TbMsgDelayNode.java | 69 ++++-- .../delay/TbMsgDelayNodeConfiguration.java | 11 +- .../rule/engine/delay/TbMsgDelayNodeTest.java | 234 ++++++++++++------ 3 files changed, 219 insertions(+), 95 deletions(-) diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/delay/TbMsgDelayNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/delay/TbMsgDelayNode.java index b60a46ee0a..1ade5efb82 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/delay/TbMsgDelayNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/delay/TbMsgDelayNode.java @@ -15,8 +15,9 @@ */ package org.thingsboard.rule.engine.delay; +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.node.ObjectNode; import lombok.extern.slf4j.Slf4j; -import org.apache.commons.lang3.math.NumberUtils; import org.thingsboard.rule.engine.api.RuleNode; import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbNode; @@ -26,10 +27,12 @@ import org.thingsboard.rule.engine.api.util.TbNodeUtils; 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.common.msg.TbMsgMetaData; import java.util.HashMap; +import java.util.List; import java.util.Map; import java.util.UUID; import java.util.concurrent.TimeUnit; @@ -38,6 +41,7 @@ import java.util.concurrent.TimeUnit; @RuleNode( type = ComponentType.ACTION, name = "delay (deprecated)", + version = 1, configClazz = TbMsgDelayNodeConfiguration.class, nodeDescription = "Delays incoming message (deprecated)", nodeDetails = "Delays messages for a configurable period. " + @@ -45,7 +49,7 @@ import java.util.concurrent.TimeUnit; "Deprecated because the acknowledged message still stays in memory (to be delayed) and this " + "does not guarantee that message will be processed even if the \"retry failures and timeouts\" processing strategy will be chosen.", icon = "pause", - uiResources = {"static/rulenode/rulenode-core-config.js"}, + uiResources = {""}, configDirective = "tbActionNodeMsgDelayConfig" ) public class TbMsgDelayNode implements TbNode { @@ -89,25 +93,58 @@ public class TbMsgDelayNode implements TbNode { } private long getDelay(TbMsg msg) { - int periodInSeconds; - if (config.isUseMetadataPeriodInSecondsPatterns()) { - if (isParsable(msg, config.getPeriodInSecondsPattern())) { - periodInSeconds = Integer.parseInt(TbNodeUtils.processPattern(config.getPeriodInSecondsPattern(), msg)); - } else { - throw new RuntimeException("Can't parse period in seconds from metadata using pattern: " + config.getPeriodInSecondsPattern()); - } - } else { - periodInSeconds = config.getPeriodInSeconds(); + String timeUnitPattern = TbNodeUtils.processPattern(config.getTimeUnit(), msg); + String periodPattern = TbNodeUtils.processPattern(config.getPeriod(), msg); + try { + TimeUnit timeUnit = TimeUnit.valueOf(timeUnitPattern.toUpperCase()); + int period = Integer.parseInt(periodPattern); + return timeUnit.toMillis(period); + } catch (NumberFormatException e) { + throw new RuntimeException("Can't parse period value : " + periodPattern); + } catch (IllegalArgumentException e) { + throw new RuntimeException("Invalid value for period time unit : " + timeUnitPattern); } - return TimeUnit.SECONDS.toMillis(periodInSeconds); - } - - private boolean isParsable(TbMsg msg, String pattern) { - return NumberUtils.isParsable(TbNodeUtils.processPattern(pattern, msg)); } @Override public void destroy() { pendingMsgs.clear(); } + + @Override + public TbPair upgrade(int fromVersion, JsonNode oldConfiguration) throws TbNodeException { + boolean hasChanges = false; + switch (fromVersion) { + case 0: + var periodInSeconds = "periodInSeconds"; + var periodInSecondsPattern = "periodInSecondsPattern"; + var useMetadataPeriodInSecondsPatterns = "useMetadataPeriodInSecondsPatterns"; + var period = "period"; + if (oldConfiguration.has(useMetadataPeriodInSecondsPatterns)) { + var isUsedPattern = oldConfiguration.get(useMetadataPeriodInSecondsPatterns).asBoolean(); + if (isUsedPattern) { + if (!oldConfiguration.has(periodInSecondsPattern)) { + throw new TbNodeException("Property to update: '" + periodInSecondsPattern + "' does not exist in configuration."); + } + ((ObjectNode) oldConfiguration).set(period, oldConfiguration.get(periodInSecondsPattern)); + } else { + if (!oldConfiguration.has(periodInSeconds)) { + throw new TbNodeException("Property to update: '" + periodInSeconds + "' does not exist in configuration."); + } + ((ObjectNode) oldConfiguration).set(period, oldConfiguration.get(periodInSeconds)); + } + ((ObjectNode) oldConfiguration).remove(List.of(periodInSeconds, periodInSecondsPattern, useMetadataPeriodInSecondsPatterns)); + hasChanges = true; + } + var timeUnit = "timeUnit"; + if (!oldConfiguration.has(timeUnit)) { + ((ObjectNode) oldConfiguration).put(timeUnit, TimeUnit.SECONDS.name()); + hasChanges = true; + } + break; + default: + break; + } + return new TbPair<>(hasChanges, oldConfiguration); + } } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/delay/TbMsgDelayNodeConfiguration.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/delay/TbMsgDelayNodeConfiguration.java index f35552de18..b1ec01bb10 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/delay/TbMsgDelayNodeConfiguration.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/delay/TbMsgDelayNodeConfiguration.java @@ -18,20 +18,21 @@ package org.thingsboard.rule.engine.delay; import lombok.Data; import org.thingsboard.rule.engine.api.NodeConfiguration; +import java.util.concurrent.TimeUnit; + @Data public class TbMsgDelayNodeConfiguration implements NodeConfiguration { - private int periodInSeconds; + private String period; + private String timeUnit; private int maxPendingMsgs; - private String periodInSecondsPattern; - private boolean useMetadataPeriodInSecondsPatterns; @Override public TbMsgDelayNodeConfiguration defaultConfiguration() { TbMsgDelayNodeConfiguration configuration = new TbMsgDelayNodeConfiguration(); - configuration.setPeriodInSeconds(60); + configuration.setPeriod("60"); + configuration.setTimeUnit(TimeUnit.SECONDS.name()); configuration.setMaxPendingMsgs(1000); - configuration.setUseMetadataPeriodInSecondsPatterns(false); return configuration; } } diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/delay/TbMsgDelayNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/delay/TbMsgDelayNodeTest.java index 18d3ffc2c5..751faee18b 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/delay/TbMsgDelayNodeTest.java +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/delay/TbMsgDelayNodeTest.java @@ -26,7 +26,9 @@ import org.mockito.Mock; import org.mockito.junit.jupiter.MockitoExtension; import org.springframework.test.util.ReflectionTestUtils; import org.thingsboard.common.util.JacksonUtil; +import org.thingsboard.rule.engine.AbstractRuleNodeUpgradeTest; import org.thingsboard.rule.engine.api.TbContext; +import org.thingsboard.rule.engine.api.TbNode; import org.thingsboard.rule.engine.api.TbNodeConfiguration; import org.thingsboard.rule.engine.api.TbNodeException; import org.thingsboard.server.common.data.id.DeviceId; @@ -36,9 +38,12 @@ import org.thingsboard.server.common.data.msg.TbNodeConnectionType; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgMetaData; +import java.util.ArrayList; import java.util.HashMap; +import java.util.List; import java.util.Map; import java.util.UUID; +import java.util.concurrent.TimeUnit; import java.util.stream.Stream; import static org.assertj.core.api.Assertions.assertThat; @@ -46,15 +51,15 @@ import static org.assertj.core.api.Assertions.assertThatNoException; import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.eq; -import static org.mockito.ArgumentMatchers.isNull; import static org.mockito.BDDMockito.given; +import static org.mockito.BDDMockito.never; import static org.mockito.BDDMockito.spy; import static org.mockito.BDDMockito.then; import static org.mockito.BDDMockito.times; import static org.mockito.BDDMockito.willAnswer; @ExtendWith(MockitoExtension.class) -public class TbMsgDelayNodeTest { +public class TbMsgDelayNodeTest extends AbstractRuleNodeUpgradeTest { private final DeviceId DEVICE_ID = new DeviceId(UUID.fromString("20107cf0-1c5e-4ac4-8131-7c466c955a7c")); @@ -72,10 +77,9 @@ public class TbMsgDelayNodeTest { @Test public void verifyDefaultConfig() { - assertThat(config.getPeriodInSeconds()).isEqualTo(60); + assertThat(config.getPeriod()).isEqualTo("60"); assertThat(config.getMaxPendingMsgs()).isEqualTo(1000); - assertThat(config.isUseMetadataPeriodInSecondsPatterns()).isFalse(); - assertThat(config.getPeriodInSecondsPattern()).isNull(); + assertThat(config.getTimeUnit()).isEqualTo(TimeUnit.SECONDS.name()); } @Test @@ -85,121 +89,137 @@ public class TbMsgDelayNodeTest { @ParameterizedTest @MethodSource - public void givenPeriodInSecondsPattern_whenOnMsg_thenTellSelfTickMsgAndEnqueueForTellNext( - String periodInSecondsPattern, TbMsgMetaData metaData, String data, long expectedDelay) throws TbNodeException { - config.setUseMetadataPeriodInSecondsPatterns(true); - config.setPeriodInSecondsPattern(periodInSecondsPattern); - node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))); + public void givenPeriodValueAndPeriodTimeUnitPatterns_whenOnMsg_thenTellSelfTickMsgAndEnqueueForTellNext( + String periodPattern, String timeUnitPattern, TbMsgMetaData metaData, String data, long expectedDelay) throws TbNodeException { + config.setPeriod(periodPattern); + config.setTimeUnit(timeUnitPattern); - TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, metaData, data); - RuleNodeId ruleNodeId = new RuleNodeId(UUID.fromString("5236e9b9-1e29-4b95-b219-7043ff8f0414")); - TbMsg tickMsg = TbMsg.newMsg(TbMsgType.DELAY_TIMEOUT_SELF_MSG, ruleNodeId, TbMsgMetaData.EMPTY, msg.getId().toString()); + node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))); + var ruleNodeId = new RuleNodeId(UUID.fromString("e8172ef8-bf91-4821-b9f5-ccd7b865e418")); given(ctxMock.getSelfId()).willReturn(ruleNodeId); - given(ctxMock.newMsg(any(), any(TbMsgType.class), any(), any(), any(), any())).willReturn(tickMsg); willAnswer(invocation -> { node.onMsg(ctxMock, invocation.getArgument(0)); return null; - }).given(ctxMock).tellSelf(tickMsg, expectedDelay); + }).given(ctxMock).tellSelf(any(TbMsg.class), any(Long.class)); - node.onMsg(ctxMock, msg); + List incomingMsgs = new ArrayList<>(); + List 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); + } - then(ctxMock).should().newMsg(null, TbMsgType.DELAY_TIMEOUT_SELF_MSG, ruleNodeId, null, TbMsgMetaData.EMPTY, msg.getId().toString()); - then(ctxMock).should().tellSelf(tickMsg, expectedDelay); - then(ctxMock).should().ack(msg); - then(node).should().onMsg(ctxMock, tickMsg); - ArgumentCaptor actualMsg = ArgumentCaptor.forClass(TbMsg.class); - then(ctxMock).should().enqueueForTellNext(actualMsg.capture(), eq(TbNodeConnectionType.SUCCESS)); - assertThat(actualMsg.getValue()).usingRecursiveComparison().ignoringFields("id", "ts").isEqualTo(msg); + incomingMsgs.forEach(msg -> node.onMsg(ctxMock, msg)); + + 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); + }); + 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 givenPeriodInSecondsPattern_whenOnMsg_thenTellSelfTickMsgAndEnqueueForTellNext() { + private static Stream givenPeriodValueAndPeriodTimeUnitPatterns_whenOnMsg_thenTellSelfTickMsgAndEnqueueForTellNext() { return Stream.of( - Arguments.of("${md-period-in-seconds}", new TbMsgMetaData(Map.of("md-period-in-seconds", "10")), TbMsg.EMPTY_JSON_OBJECT, 10000L), - Arguments.of("$[msg-period-in-seconds]", TbMsgMetaData.EMPTY, "{\"msg-period-in-seconds\":5}", 5000L) + Arguments.of("1", "HOURS", TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT, TimeUnit.HOURS.toMillis(1L)), + Arguments.of("${md-period}", "${md-time-unit}", + new TbMsgMetaData(Map.of( + "md-period", "5", + "md-time-unit", "MINUTES" + )), TbMsg.EMPTY_JSON_OBJECT, TimeUnit.MINUTES.toMillis(5L)), + Arguments.of("$[msg-period]", "$[msg-time-unit]", TbMsgMetaData.EMPTY, + "{\"msg-period\":10,\"msg-time-unit\":\"SECONDS\"}", TimeUnit.SECONDS.toMillis(10L)) ); } @Test - public void givenPeriodInSecondsPatternIsUnparsable_whenOnMsg_thenThrowsException() throws TbNodeException { - config.setUseMetadataPeriodInSecondsPatterns(true); - config.setPeriodInSecondsPattern("$[msg-period-in-seconds]"); + public void givenPeriodIsUnparsable_whenOnMsg_thenThrowsException() throws TbNodeException { + config.setPeriod("five"); node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))); - TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, TbMsgMetaData.EMPTY, "{\"msg-period-in-seconds\":\"five\"}"); - RuleNodeId ruleNodeId = new RuleNodeId(UUID.fromString("5236e9b9-1e29-4b95-b219-7043ff8f0414")); - TbMsg tickMsg = TbMsg.newMsg(TbMsgType.DELAY_TIMEOUT_SELF_MSG, ruleNodeId, TbMsgMetaData.EMPTY, msg.getId().toString()); + 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)) .isInstanceOf(RuntimeException.class) - .hasMessage("Can't parse period in seconds from metadata using pattern: $[msg-period-in-seconds]"); + .hasMessage("Can't parse period value : five"); } @Test - public void givenMaxLimitOfPendingMsgsReached_whenOnMsg_thenTellFailure() throws TbNodeException { - int maxPendingMsgs = 5; - config.setMaxPendingMsgs(maxPendingMsgs); - + public void givenInvalidTimeUnit_whenOnMsg_thenThrowsException() throws TbNodeException { + config.setTimeUnit("sec"); node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))); - RuleNodeId ruleNodeId = new RuleNodeId(UUID.fromString("d1440f09-ca81-41f3-b67e-1495aee87dc6")); - String msgId = "e38c87c5-8916-4bb0-b448-b8d08ad639df"; - TbMsg tickMsg = TbMsg.newMsg(TbMsgType.DELAY_TIMEOUT_SELF_MSG, ruleNodeId, TbMsgMetaData.EMPTY, msgId); + 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); - for (int i = 0; i < 6; i++) { - TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, TbMsgMetaData.EMPTY, "msg : " + (i + 1)); - node.onMsg(ctxMock, msg); - } - - then(ctxMock).should(times(maxPendingMsgs)).newMsg(isNull(), eq(TbMsgType.DELAY_TIMEOUT_SELF_MSG), eq(ruleNodeId), isNull(), eq(TbMsgMetaData.EMPTY), any()); - then(ctxMock).should(times(maxPendingMsgs)).tellSelf(tickMsg, 60000L); - then(ctxMock).should(times(maxPendingMsgs)).ack(any()); - ArgumentCaptor actualMsg = ArgumentCaptor.forClass(TbMsg.class); - ArgumentCaptor throwable = ArgumentCaptor.forClass(Throwable.class); - then(ctxMock).should().tellFailure(actualMsg.capture(), throwable.capture()); - TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, TbMsgMetaData.EMPTY, "msg : 6"); - assertThat(actualMsg.getValue()).usingRecursiveComparison().ignoringFields("id", "ts").isEqualTo(msg); - assertThat(throwable.getValue()).isInstanceOf(RuntimeException.class).hasMessage("Max limit of pending messages reached!"); + assertThatThrownBy(() -> node.onMsg(ctxMock, msg)) + .isInstanceOf(RuntimeException.class) + .hasMessage("Invalid value for period time unit : sec"); } @Test - public void givenNumberOfMsgsMoreThenMaxPendingMsgs_whenOnMsg_thenTellSelfTickMsgAndEnqueueForTellNext() throws TbNodeException { - config.setMaxPendingMsgs(3); - config.setPeriodInSeconds(1); + public void givenMaxLimitOfPendingMsgsReached_whenOnMsg_thenTellFailure() throws TbNodeException { + int maxPendingMsgs = 5; + config.setMaxPendingMsgs(maxPendingMsgs); node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))); - RuleNodeId ruleNodeId = new RuleNodeId(UUID.fromString("e8172ef8-bf91-4821-b9f5-ccd7b865e418")); - TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT); - TbMsg tickMsg = TbMsg.newMsg(TbMsgType.DELAY_TIMEOUT_SELF_MSG, ruleNodeId, TbMsgMetaData.EMPTY, msg.getId().toString()); - + RuleNodeId ruleNodeId = new RuleNodeId(UUID.fromString("d1440f09-ca81-41f3-b67e-1495aee87dc6")); given(ctxMock.getSelfId()).willReturn(ruleNodeId); - given(ctxMock.newMsg(any(), any(TbMsgType.class), any(), any(), any(), any())).willReturn(tickMsg); - willAnswer(invocation -> { - node.onMsg(ctxMock, invocation.getArgument(0)); - return null; - }).given(ctxMock).tellSelf(tickMsg, 1000L); - for (int i = 0; i < 9; i++) { - node.onMsg(ctxMock, msg); + List incomingMsgs = new ArrayList<>(); + List 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)); - then(ctxMock).should(times(9)).newMsg(null, TbMsgType.DELAY_TIMEOUT_SELF_MSG, ruleNodeId, null, TbMsgMetaData.EMPTY, msg.getId().toString()); - then(ctxMock).should(times(9)).tellSelf(tickMsg, 1000L); - then(ctxMock).should(times(9)).ack(msg); - then(node).should(times(9)).onMsg(ctxMock, tickMsg); - then(ctxMock).should(times(9)).enqueueForTellNext(any(TbMsg.class), eq(TbNodeConnectionType.SUCCESS)); + 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 = ArgumentCaptor.forClass(Throwable.class); + then(ctxMock).should().tellFailure(eq(lastMsg), throwable.capture()); + assertThat(throwable.getValue()).isInstanceOf(RuntimeException.class).hasMessage("Max limit of pending messages reached!"); } @Test public void verifyDestroyMethod() { - TbMsg 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(); pendingMsgs.put(UUID.fromString("321f0301-9bed-4e7d-b92f-a978f53ec5d6"), msg); ReflectionTestUtils.setField(node, "pendingMsgs", pendingMsgs); @@ -210,4 +230,70 @@ public class TbMsgDelayNodeTest { assertThat(actualPendingMsgs).isEmpty(); } + + private static Stream givenFromVersionAndConfig_whenUpgrade_thenVerifyHasChangesAndConfig() { + return Stream.of( + // config for version 1 with upgrade from version 0 + 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) + Arguments.of(0, + """ + { + "periodInSeconds": 60, + "maxPendingMsgs": 1000, + "periodInSecondsPattern": "${period-pattern}", + "useMetadataPeriodInSecondsPatterns": true + } + """, + true, + """ + { + "period": "${period-pattern}", + "timeUnit": "SECONDS", + "maxPendingMsgs": 1000 + } + """ + ), + // config for version 1 with upgrade from version 0 (hasChanges is false) + Arguments.of(0, + """ + { + "period": "${period-pattern}", + "timeUnit": "SECONDS", + "maxPendingMsgs": 1000 + } + """, + false, + """ + { + "period": "${period-pattern}", + "timeUnit": "SECONDS", + "maxPendingMsgs": 1000 + } + """ + ) + ); + } + + @Override + protected TbNode getTestNode() { + return node; + } } From 291849f81b09a31815e4533b4dec614f10f23b18 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Wed, 3 Jul 2024 16:36:39 +0300 Subject: [PATCH 03/12] fixed ui resources --- .../java/org/thingsboard/rule/engine/delay/TbMsgDelayNode.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/delay/TbMsgDelayNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/delay/TbMsgDelayNode.java index 1ade5efb82..3f1cd88571 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/delay/TbMsgDelayNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/delay/TbMsgDelayNode.java @@ -49,7 +49,7 @@ import java.util.concurrent.TimeUnit; "Deprecated because the acknowledged message still stays in memory (to be delayed) and this " + "does not guarantee that message will be processed even if the \"retry failures and timeouts\" processing strategy will be chosen.", icon = "pause", - uiResources = {""}, + uiResources = {"static/rulenode/rulenode-core-config.js"}, configDirective = "tbActionNodeMsgDelayConfig" ) public class TbMsgDelayNode implements TbNode { From a34088fcfbece13baccba1871e67e648988f2fc1 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Thu, 15 Aug 2024 10:17:47 +0300 Subject: [PATCH 04/12] added verification on init --- .../rule/engine/delay/TbMsgDelayNode.java | 66 +++-- .../rule/engine/delay/TbMsgDelayNodeTest.java | 249 ++++++++++-------- 2 files changed, 179 insertions(+), 136 deletions(-) diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/delay/TbMsgDelayNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/delay/TbMsgDelayNode.java index 3f1cd88571..be569162c4 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/delay/TbMsgDelayNode.java +++ b/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.TbMsgMetaData; -import java.util.HashMap; import java.util.List; -import java.util.Map; import java.util.UUID; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ConcurrentMap; import java.util.concurrent.TimeUnit; @Slf4j @@ -54,13 +54,17 @@ import java.util.concurrent.TimeUnit; ) public class TbMsgDelayNode implements TbNode { + private final List 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 Map pendingMsgs; + private ConcurrentMap pendingMsgs; @Override public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException { this.config = TbNodeUtils.convert(configuration, TbMsgDelayNodeConfiguration.class); - this.pendingMsgs = new HashMap<>(); + validateConfig(); + this.pendingMsgs = new ConcurrentHashMap<>(); } @Override @@ -71,7 +75,7 @@ public class TbMsgDelayNode implements TbNode { ctx.enqueueForTellNext( TbMsg.newMsg( pendingMsg.getQueueName(), - pendingMsg.getType(), + pendingMsg.getInternalType(), pendingMsg.getOriginator(), pendingMsg.getCustomerId(), pendingMsg.getMetaData(), @@ -97,12 +101,28 @@ public class TbMsgDelayNode implements TbNode { String periodPattern = TbNodeUtils.processPattern(config.getPeriod(), msg); try { 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); return timeUnit.toMillis(period); } catch (NumberFormatException e) { - throw new RuntimeException("Can't parse period value : " + periodPattern); + throw new NumberFormatException("Can't parse period value : " + periodPattern); } 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 period = "period"; if (oldConfiguration.has(useMetadataPeriodInSecondsPatterns)) { - var isUsedPattern = oldConfiguration.get(useMetadataPeriodInSecondsPatterns).asBoolean(); - if (isUsedPattern) { - if (!oldConfiguration.has(periodInSecondsPattern)) { - throw new TbNodeException("Property to update: '" + periodInSecondsPattern + "' does not exist in configuration."); - } - ((ObjectNode) oldConfiguration).set(period, oldConfiguration.get(periodInSecondsPattern)); - } else { - if (!oldConfiguration.has(periodInSeconds)) { - throw new TbNodeException("Property to update: '" + periodInSeconds + "' does not exist in configuration."); - } - ((ObjectNode) oldConfiguration).set(period, oldConfiguration.get(periodInSeconds)); - } - ((ObjectNode) oldConfiguration).remove(List.of(periodInSeconds, periodInSecondsPattern, useMetadataPeriodInSecondsPatterns)); - hasChanges = true; + var isUsedPattern = oldConfiguration.get(useMetadataPeriodInSecondsPatterns).booleanValue(); + if (isUsedPattern) { + if (!oldConfiguration.has(periodInSecondsPattern)) { + throw new TbNodeException("Property to update: '" + periodInSecondsPattern + "' does not exist in configuration."); + } + ((ObjectNode) oldConfiguration).set(period, oldConfiguration.get(periodInSecondsPattern)); + } else { + if (!oldConfiguration.has(periodInSeconds)) { + throw new TbNodeException("Property to update: '" + periodInSeconds + "' does not exist in configuration."); + } + ((ObjectNode) oldConfiguration).put(period, oldConfiguration.get(periodInSeconds).asText()); + } + hasChanges = true; + } + if (!oldConfiguration.has(period)) { + ((ObjectNode) oldConfiguration).put(period, "60"); + hasChanges = true; } var timeUnit = "timeUnit"; if (!oldConfiguration.has(timeUnit)) { ((ObjectNode) oldConfiguration).put(timeUnit, TimeUnit.SECONDS.name()); hasChanges = true; } + ((ObjectNode) oldConfiguration).remove(List.of(periodInSeconds, periodInSecondsPattern, useMetadataPeriodInSecondsPatterns)); break; default: break; diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/delay/TbMsgDelayNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/delay/TbMsgDelayNodeTest.java index 751faee18b..63de8b505b 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/delay/TbMsgDelayNodeTest.java +++ b/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.params.ParameterizedTest; 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.ValueSource; import org.mockito.ArgumentCaptor; import org.mockito.Mock; 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.TbMsgMetaData; -import java.util.ArrayList; -import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.UUID; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.TimeUnit; 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.eq; import static org.mockito.BDDMockito.given; -import static org.mockito.BDDMockito.never; import static org.mockito.BDDMockito.spy; import static org.mockito.BDDMockito.then; -import static org.mockito.BDDMockito.times; import static org.mockito.BDDMockito.willAnswer; @ExtendWith(MockitoExtension.class) public class TbMsgDelayNodeTest extends AbstractRuleNodeUpgradeTest { 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 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 TbMsgDelayNodeConfiguration config; @@ -87,6 +90,40 @@ public class TbMsgDelayNodeTest extends AbstractRuleNodeUpgradeTest { 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 @MethodSource public void givenPeriodValueAndPeriodTimeUnitPatterns_whenOnMsg_thenTellSelfTickMsgAndEnqueueForTellNext( @@ -96,41 +133,23 @@ public class TbMsgDelayNodeTest extends AbstractRuleNodeUpgradeTest { node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))); - var ruleNodeId = new RuleNodeId(UUID.fromString("e8172ef8-bf91-4821-b9f5-ccd7b865e418")); - given(ctxMock.getSelfId()).willReturn(ruleNodeId); + var msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, metaData, data); + 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 -> { node.onMsg(ctxMock, invocation.getArgument(0)); return null; }).given(ctxMock).tellSelf(any(TbMsg.class), any(Long.class)); - List incomingMsgs = new ArrayList<>(); - List 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); - } + node.onMsg(ctxMock, msg); - incomingMsgs.forEach(msg -> node.onMsg(ctxMock, msg)); - - 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); - }); - 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)); - } + then(ctxMock).should().tellSelf(tickMsg, expectedDelay); + then(ctxMock).should().ack(msg); + ArgumentCaptor actualMsg = ArgumentCaptor.forClass(TbMsg.class); + then(ctxMock).should().enqueueForTellNext(actualMsg.capture(), eq(TbNodeConnectionType.SUCCESS)); + assertThat(actualMsg.getValue()).usingRecursiveComparison().ignoringFields("id", "ts").isEqualTo(msg); } private static Stream 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("${md-period}", "${md-time-unit}", new TbMsgMetaData(Map.of( - "md-period", "5", - "md-time-unit", "MINUTES" + "md-period", "5", + "md-time-unit", "MINUTES" )), TbMsg.EMPTY_JSON_OBJECT, TimeUnit.MINUTES.toMillis(5L)), Arguments.of("$[msg-period]", "$[msg-time-unit]", TbMsgMetaData.EMPTY, "{\"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 public void givenPeriodIsUnparsable_whenOnMsg_thenThrowsException() throws TbNodeException { config.setPeriod("five"); 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 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)) - .isInstanceOf(RuntimeException.class) + .isInstanceOf(NumberFormatException.class) .hasMessage("Can't parse period value : five"); } @@ -169,58 +198,30 @@ public class TbMsgDelayNodeTest extends AbstractRuleNodeUpgradeTest { 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 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)) - .isInstanceOf(RuntimeException.class) + .isInstanceOf(IllegalArgumentException.class) .hasMessage("Invalid value for period time unit : sec"); } @Test public void givenMaxLimitOfPendingMsgsReached_whenOnMsg_thenTellFailure() throws TbNodeException { - int maxPendingMsgs = 5; - config.setMaxPendingMsgs(maxPendingMsgs); + config.setMaxPendingMsgs(1); node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))); - - RuleNodeId ruleNodeId = new RuleNodeId(UUID.fromString("d1440f09-ca81-41f3-b67e-1495aee87dc6")); - given(ctxMock.getSelfId()).willReturn(ruleNodeId); - - List incomingMsgs = new ArrayList<>(); - List 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)); + var msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT); + for (int i = 0; i < 2; i++) { + node.onMsg(ctxMock, msg); } - 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 = 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!"); } @Test public void verifyDestroyMethod() { var msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT); - var pendingMsgs = new HashMap(); + var pendingMsgs = new ConcurrentHashMap<>(); pendingMsgs.put(UUID.fromString("321f0301-9bed-4e7d-b92f-a978f53ec5d6"), msg); ReflectionTestUtils.setField(node, "pendingMsgs", pendingMsgs); var actualPendingMsgs = (Map) ReflectionTestUtils.getField(node, "pendingMsgs"); @@ -233,61 +234,79 @@ public class TbMsgDelayNodeTest extends AbstractRuleNodeUpgradeTest { private static Stream givenFromVersionAndConfig_whenUpgrade_thenVerifyHasChangesAndConfig() { 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, """ - { - "periodInSeconds": 60, - "maxPendingMsgs": 1000, - "periodInSecondsPattern": null, - "useMetadataPeriodInSecondsPatterns": false - } - """, + { + "periodInSeconds": 13, + "maxPendingMsgs": 1000, + "periodInSecondsPattern": "17" + } + """, true, """ - { - "period": 60, - "timeUnit": "SECONDS", - "maxPendingMsgs": 1000 - } + { + "period": "60", + "timeUnit": "SECONDS", + "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) Arguments.of(0, """ - { - "periodInSeconds": 60, - "maxPendingMsgs": 1000, - "periodInSecondsPattern": "${period-pattern}", - "useMetadataPeriodInSecondsPatterns": true - } - """, + { + "periodInSeconds": 60, + "maxPendingMsgs": 1000, + "periodInSecondsPattern": "${period-pattern}", + "useMetadataPeriodInSecondsPatterns": true + } + """, true, """ - { - "period": "${period-pattern}", - "timeUnit": "SECONDS", - "maxPendingMsgs": 1000 - } - """ + { + "period": "${period-pattern}", + "timeUnit": "SECONDS", + "maxPendingMsgs": 1000 + } + """ ), // config for version 1 with upgrade from version 0 (hasChanges is false) Arguments.of(0, """ - { - "period": "${period-pattern}", - "timeUnit": "SECONDS", - "maxPendingMsgs": 1000 - } - """, + { + "period": "${period-pattern}", + "timeUnit": "SECONDS", + "maxPendingMsgs": 1000 + } + """, false, """ - { - "period": "${period-pattern}", - "timeUnit": "SECONDS", - "maxPendingMsgs": 1000 - } - """ + { + "period": "${period-pattern}", + "timeUnit": "SECONDS", + "maxPendingMsgs": 1000 + } + """ ) ); } From 8c681508d51efc811114f41eff4c3cc69849b8f0 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Thu, 22 Aug 2024 15:05:59 +0300 Subject: [PATCH 05/12] changed list to set --- .../thingsboard/rule/engine/delay/TbMsgDelayNode.java | 6 ++++-- .../rule/engine/delay/TbMsgDelayNodeTest.java | 9 +++++---- 2 files changed, 9 insertions(+), 6 deletions(-) diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/delay/TbMsgDelayNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/delay/TbMsgDelayNode.java index be569162c4..a336c2095f 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/delay/TbMsgDelayNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/delay/TbMsgDelayNode.java @@ -31,7 +31,9 @@ import org.thingsboard.server.common.data.util.TbPair; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgMetaData; +import java.util.EnumSet; import java.util.List; +import java.util.Set; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; @@ -54,8 +56,8 @@ import java.util.concurrent.TimeUnit; ) public class TbMsgDelayNode implements TbNode { - private final List 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 static final Set supportedTimeUnits = EnumSet.of(TimeUnit.SECONDS, TimeUnit.MINUTES, TimeUnit.HOURS); + private static final String supportedTimeUnitsStr = String.join(",", TimeUnit.SECONDS.name(), TimeUnit.MINUTES.name(), TimeUnit.HOURS.name()); private TbMsgDelayNodeConfiguration config; private ConcurrentMap pendingMsgs; diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/delay/TbMsgDelayNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/delay/TbMsgDelayNodeTest.java index 63de8b505b..06565465d9 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/delay/TbMsgDelayNodeTest.java +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/delay/TbMsgDelayNodeTest.java @@ -40,8 +40,9 @@ import org.thingsboard.server.common.data.msg.TbNodeConnectionType; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgMetaData; -import java.util.List; +import java.util.EnumSet; import java.util.Map; +import java.util.Set; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.TimeUnit; @@ -63,7 +64,7 @@ public class TbMsgDelayNodeTest extends AbstractRuleNodeUpgradeTest { 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 supportedTimeUnits = List.of(TimeUnit.SECONDS, TimeUnit.MINUTES, TimeUnit.HOURS); + private final Set supportedTimeUnits = EnumSet.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; @@ -92,8 +93,8 @@ public class TbMsgDelayNodeTest extends AbstractRuleNodeUpgradeTest { @ParameterizedTest @ValueSource(ints = {-1, 0, 5000000}) - public void givenInvalidMaxPendingMsgsValue_whenInit_thenThrowsException() { - config.setMaxPendingMsgs(-1); + public void givenInvalidMaxPendingMsgsValue_whenInit_thenThrowsException(int maxPendingMsgs) { + config.setMaxPendingMsgs(maxPendingMsgs); assertThatThrownBy(() -> node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)))) .isInstanceOf(TbNodeException.class) From aac427d6bee41711006ff19f1bda5819f113ee0b Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Fri, 23 Aug 2024 12:34:43 +0300 Subject: [PATCH 06/12] added one more verification in upgrade --- .../rule/engine/delay/TbMsgDelayNode.java | 3 +++ .../rule/engine/delay/TbMsgDelayNodeTest.java | 19 ++++++++++++++++++- 2 files changed, 21 insertions(+), 1 deletion(-) diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/delay/TbMsgDelayNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/delay/TbMsgDelayNode.java index a336c2095f..a3e0f6becc 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/delay/TbMsgDelayNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/delay/TbMsgDelayNode.java @@ -156,6 +156,9 @@ public class TbMsgDelayNode implements TbNode { ((ObjectNode) oldConfiguration).put(period, oldConfiguration.get(periodInSeconds).asText()); } hasChanges = true; + } else if (oldConfiguration.has(periodInSeconds)) { + ((ObjectNode) oldConfiguration).put(period, oldConfiguration.get(periodInSeconds).asText()); + hasChanges = true; } if (!oldConfiguration.has(period)) { ((ObjectNode) oldConfiguration).put(period, "60"); diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/delay/TbMsgDelayNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/delay/TbMsgDelayNodeTest.java index 06565465d9..70511b03b9 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/delay/TbMsgDelayNodeTest.java +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/delay/TbMsgDelayNodeTest.java @@ -235,7 +235,7 @@ public class TbMsgDelayNodeTest extends AbstractRuleNodeUpgradeTest { private static Stream givenFromVersionAndConfig_whenUpgrade_thenVerifyHasChangesAndConfig() { return Stream.of( - // config for version 1 with upgrade from version 0 (useMetadataPeriodInSecondsPatterns does not exist) + // config for version 1 with upgrade from version 0 (useMetadataPeriodInSecondsPatterns does not exist and periodInSeconds exists) Arguments.of(0, """ { @@ -245,6 +245,23 @@ public class TbMsgDelayNodeTest extends AbstractRuleNodeUpgradeTest { } """, true, + """ + { + "period": "13", + "timeUnit": "SECONDS", + "maxPendingMsgs": 1000 + } + """ + ), + // config for version 1 with upgrade from version 0 (useMetadataPeriodInSecondsPatterns and periodInSeconds do not exist) + Arguments.of(0, + """ + { + "maxPendingMsgs": 1000, + "periodInSecondsPattern": "17" + } + """, + true, """ { "period": "60", From 53179f8007b98915b871122be51f0185a7301540 Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Tue, 10 Sep 2024 13:15:28 +0200 Subject: [PATCH 07/12] removed wrong CONNECTION close header --- .../main/java/org/thingsboard/rule/engine/rest/TbHttpClient.java | 1 - 1 file changed, 1 deletion(-) diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbHttpClient.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbHttpClient.java index f8ec81abcf..d7476bf08e 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbHttpClient.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbHttpClient.java @@ -135,7 +135,6 @@ public class TbHttpClient { this.webClient = WebClient.builder() .clientConnector(new ReactorClientHttpConnector(httpClient)) - .defaultHeader(HttpHeaders.CONNECTION, "close") //In previous realization this header was present! (Added for hotfix "Connection reset") .codecs(configurer -> configurer.defaultCodecs().maxInMemorySize( (config.getMaxInMemoryBufferSizeInKb() > 0 ? config.getMaxInMemoryBufferSizeInKb() : 256) * 1024)) .build(); From 4b697482cac4a48931b81c1e9136ef57fee8bf28 Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Tue, 10 Sep 2024 18:19:04 +0200 Subject: [PATCH 08/12] added ability to configure maxConnections in TbHttpClient using env TB_POOL_MAX_CONNECTIONS --- .../rule/engine/rest/TbHttpClient.java | 20 ++++++++++++++++++- 1 file changed, 19 insertions(+), 1 deletion(-) diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbHttpClient.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbHttpClient.java index d7476bf08e..b0b41fba4b 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbHttpClient.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbHttpClient.java @@ -42,6 +42,7 @@ import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgMetaData; import reactor.netty.http.client.HttpClient; +import reactor.netty.resources.ConnectionProvider; import reactor.netty.transport.ProxyProvider; import javax.net.ssl.SSLException; @@ -95,7 +96,12 @@ public class TbHttpClient { semaphore = new Semaphore(config.getMaxParallelRequestsCount()); } - HttpClient httpClient = HttpClient.create() + ConnectionProvider connectionProvider = ConnectionProvider + .builder("http") + .maxConnections(getPoolMaxConnections()) + .build(); + + HttpClient httpClient = HttpClient.create(connectionProvider) .runOn(getSharedOrCreateEventLoopGroup(eventLoopGroupShared)) .doOnConnected(c -> c.addHandlerLast(new ReadTimeoutHandler(config.getReadTimeoutMs(), TimeUnit.MILLISECONDS))); @@ -143,6 +149,18 @@ public class TbHttpClient { } } + private int getPoolMaxConnections() { + String poolMaxConnectionsEnv = System.getenv("TB_POOL_MAX_CONNECTIONS"); + + int poolMaxConnections; + if (poolMaxConnectionsEnv != null) { + poolMaxConnections = Integer.parseInt(poolMaxConnectionsEnv); + } else { + poolMaxConnections = ConnectionProvider.DEFAULT_POOL_MAX_CONNECTIONS; + } + return poolMaxConnections; + } + private void validateMaxInMemoryBufferSize(TbRestApiCallNodeConfiguration config) throws TbNodeException { int systemMaxInMemoryBufferSizeInKb = 25000; try { From 4431a32659082ddd1b2cefb3050e6983c9df6d66 Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Wed, 11 Sep 2024 12:09:11 +0200 Subject: [PATCH 09/12] minor refactoring due to comments --- .../java/org/thingsboard/rule/engine/rest/TbHttpClient.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbHttpClient.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbHttpClient.java index b0b41fba4b..c7f0109271 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbHttpClient.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbHttpClient.java @@ -150,7 +150,7 @@ public class TbHttpClient { } private int getPoolMaxConnections() { - String poolMaxConnectionsEnv = System.getenv("TB_POOL_MAX_CONNECTIONS"); + String poolMaxConnectionsEnv = System.getenv("TB_HTTP_POOL_MAX_CONNECTIONS"); int poolMaxConnections; if (poolMaxConnectionsEnv != null) { From 54e86b27a808e265173587352d4c7af6d878bdca Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Wed, 11 Sep 2024 15:26:06 +0200 Subject: [PATCH 10/12] minor refactoring due to comments --- .../java/org/thingsboard/rule/engine/rest/TbHttpClient.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbHttpClient.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbHttpClient.java index c7f0109271..566430bd8d 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbHttpClient.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbHttpClient.java @@ -97,7 +97,7 @@ public class TbHttpClient { } ConnectionProvider connectionProvider = ConnectionProvider - .builder("http") + .builder("rule-engine-http-client") .maxConnections(getPoolMaxConnections()) .build(); @@ -150,7 +150,7 @@ public class TbHttpClient { } private int getPoolMaxConnections() { - String poolMaxConnectionsEnv = System.getenv("TB_HTTP_POOL_MAX_CONNECTIONS"); + String poolMaxConnectionsEnv = System.getenv("TB_RE_HTTP_CLIENT_POOL_MAX_CONNECTIONS"); int poolMaxConnections; if (poolMaxConnectionsEnv != null) { From 3ba6ac45e9815fba8ee59ca3916a28d7933169ae Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Thu, 12 Sep 2024 14:06:36 +0300 Subject: [PATCH 11/12] fixed error message and added replaced validation method with jakarta annotations --- .../rule/engine/delay/TbMsgDelayNode.java | 25 ++++----- .../delay/TbMsgDelayNodeConfiguration.java | 7 +++ .../rule/engine/delay/TbMsgDelayNodeTest.java | 51 +++++++++++-------- 3 files changed, 47 insertions(+), 36 deletions(-) diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/delay/TbMsgDelayNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/delay/TbMsgDelayNode.java index a3e0f6becc..267d3fae73 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/delay/TbMsgDelayNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/delay/TbMsgDelayNode.java @@ -30,6 +30,7 @@ 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.common.msg.TbMsgMetaData; +import org.thingsboard.server.dao.exception.DataValidationException; import java.util.EnumSet; import java.util.List; @@ -38,6 +39,9 @@ import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.TimeUnit; +import java.util.stream.Collectors; + +import static org.thingsboard.server.dao.service.ConstraintValidator.validateFields; @Slf4j @RuleNode( @@ -57,7 +61,7 @@ import java.util.concurrent.TimeUnit; public class TbMsgDelayNode implements TbNode { private static final Set supportedTimeUnits = EnumSet.of(TimeUnit.SECONDS, TimeUnit.MINUTES, TimeUnit.HOURS); - private static final String supportedTimeUnitsStr = String.join(",", TimeUnit.SECONDS.name(), TimeUnit.MINUTES.name(), TimeUnit.HOURS.name()); + private static final String supportedTimeUnitsStr = supportedTimeUnits.stream().map(TimeUnit::name).collect(Collectors.joining(", ")); private TbMsgDelayNodeConfiguration config; private ConcurrentMap pendingMsgs; @@ -65,7 +69,12 @@ public class TbMsgDelayNode implements TbNode { @Override public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException { this.config = TbNodeUtils.convert(configuration, TbMsgDelayNodeConfiguration.class); - validateConfig(); + String errorPrefix = "'" + ctx.getSelf().getName() + "' node configuration is invalid: "; + try { + validateFields(config, errorPrefix); + } catch (DataValidationException e) { + throw new TbNodeException(e, true); + } this.pendingMsgs = new ConcurrentHashMap<>(); } @@ -116,18 +125,6 @@ public class TbMsgDelayNode implements TbNode { } } - 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); - } - } - @Override public void destroy() { pendingMsgs.clear(); diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/delay/TbMsgDelayNodeConfiguration.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/delay/TbMsgDelayNodeConfiguration.java index b1ec01bb10..56df2c8aea 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/delay/TbMsgDelayNodeConfiguration.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/delay/TbMsgDelayNodeConfiguration.java @@ -15,6 +15,9 @@ */ package org.thingsboard.rule.engine.delay; +import jakarta.validation.constraints.Max; +import jakarta.validation.constraints.Min; +import jakarta.validation.constraints.NotNull; import lombok.Data; import org.thingsboard.rule.engine.api.NodeConfiguration; @@ -23,8 +26,12 @@ import java.util.concurrent.TimeUnit; @Data public class TbMsgDelayNodeConfiguration implements NodeConfiguration { + @NotNull private String period; + @NotNull private String timeUnit; + @Min(1) + @Max(100000) private int maxPendingMsgs; @Override diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/delay/TbMsgDelayNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/delay/TbMsgDelayNodeTest.java index 70511b03b9..23fb69d8b7 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/delay/TbMsgDelayNodeTest.java +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/delay/TbMsgDelayNodeTest.java @@ -37,6 +37,7 @@ import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.RuleNodeId; import org.thingsboard.server.common.data.msg.TbMsgType; import org.thingsboard.server.common.data.msg.TbNodeConnectionType; +import org.thingsboard.server.common.data.rule.RuleNode; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgMetaData; @@ -46,6 +47,7 @@ import java.util.Set; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.TimeUnit; +import java.util.stream.Collectors; import java.util.stream.Stream; import static org.assertj.core.api.Assertions.assertThat; @@ -65,13 +67,15 @@ public class TbMsgDelayNodeTest extends AbstractRuleNodeUpgradeTest { private final RuleNodeId RULE_NODE_ID = new RuleNodeId(UUID.fromString("1be24225-b669-4b26-ab7e-083aaa82d0a0")); private final Set supportedTimeUnits = EnumSet.of(TimeUnit.SECONDS, TimeUnit.MINUTES, TimeUnit.HOURS); - private final String supportedTimeUnitsStr = String.join(",", TimeUnit.SECONDS.name(), TimeUnit.MINUTES.name(), TimeUnit.HOURS.name()); + private final String supportedTimeUnitsStr = supportedTimeUnits.stream().map(TimeUnit::name).collect(Collectors.joining(", ")); private TbMsgDelayNode node; private TbMsgDelayNodeConfiguration config; @Mock private TbContext ctxMock; + @Mock + private RuleNode ruleNodeMock; @BeforeEach public void setUp() { @@ -88,6 +92,7 @@ public class TbMsgDelayNodeTest extends AbstractRuleNodeUpgradeTest { @Test public void givenDefaultConfig_whenInit_thenOk() { + given(ctxMock.getSelf()).willReturn(ruleNodeMock); assertThatNoException().isThrownBy(() -> node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)))); } @@ -95,34 +100,19 @@ public class TbMsgDelayNodeTest extends AbstractRuleNodeUpgradeTest { @ValueSource(ints = {-1, 0, 5000000}) public void givenInvalidMaxPendingMsgsValue_whenInit_thenThrowsException(int maxPendingMsgs) { config.setMaxPendingMsgs(maxPendingMsgs); - - 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); + verifyValidationExceptionOnInit(); } @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); + verifyValidationExceptionOnInit(); } @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); + verifyValidationExceptionOnInit(); } @ParameterizedTest @@ -131,6 +121,7 @@ public class TbMsgDelayNodeTest extends AbstractRuleNodeUpgradeTest { String periodPattern, String timeUnitPattern, TbMsgMetaData metaData, String data, long expectedDelay) throws TbNodeException { config.setPeriod(periodPattern); config.setTimeUnit(timeUnitPattern); + given(ctxMock.getSelf()).willReturn(ruleNodeMock); node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))); @@ -170,8 +161,9 @@ public class TbMsgDelayNodeTest extends AbstractRuleNodeUpgradeTest { @EnumSource(TimeUnit.class) public void givenTimeUnit_whenOnMsg_thenVerify(TimeUnit timeUnit) throws TbNodeException { config.setTimeUnit(timeUnit.name()); - node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))); + given(ctxMock.getSelf()).willReturn(ruleNodeMock); + 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)); @@ -185,8 +177,9 @@ public class TbMsgDelayNodeTest extends AbstractRuleNodeUpgradeTest { @Test public void givenPeriodIsUnparsable_whenOnMsg_thenThrowsException() throws TbNodeException { config.setPeriod("five"); - node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))); + given(ctxMock.getSelf()).willReturn(ruleNodeMock); + node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))); var msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT); assertThatThrownBy(() -> node.onMsg(ctxMock, msg)) .isInstanceOf(NumberFormatException.class) @@ -196,8 +189,9 @@ public class TbMsgDelayNodeTest extends AbstractRuleNodeUpgradeTest { @Test public void givenInvalidTimeUnit_whenOnMsg_thenThrowsException() throws TbNodeException { config.setTimeUnit("sec"); - node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))); + given(ctxMock.getSelf()).willReturn(ruleNodeMock); + node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))); var msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT); assertThatThrownBy(() -> node.onMsg(ctxMock, msg)) .isInstanceOf(IllegalArgumentException.class) @@ -207,6 +201,7 @@ public class TbMsgDelayNodeTest extends AbstractRuleNodeUpgradeTest { @Test public void givenMaxLimitOfPendingMsgsReached_whenOnMsg_thenTellFailure() throws TbNodeException { config.setMaxPendingMsgs(1); + given(ctxMock.getSelf()).willReturn(ruleNodeMock); node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))); var msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT); @@ -233,6 +228,18 @@ public class TbMsgDelayNodeTest extends AbstractRuleNodeUpgradeTest { assertThat(actualPendingMsgs).isEmpty(); } + private void verifyValidationExceptionOnInit() { + RuleNode ruleNode = new RuleNode(); + ruleNode.setName("test"); + given(ctxMock.getSelf()).willReturn(ruleNode); + String errorPrefix = "'test' node configuration is invalid: "; + assertThatThrownBy(() -> node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)))) + .isInstanceOf(TbNodeException.class) + .hasMessageContaining(errorPrefix) + .extracting(e -> ((TbNodeException) e).isUnrecoverable()) + .isEqualTo(true); + } + private static Stream givenFromVersionAndConfig_whenUpgrade_thenVerifyHasChangesAndConfig() { return Stream.of( // config for version 1 with upgrade from version 0 (useMetadataPeriodInSecondsPatterns does not exist and periodInSeconds exists) From 29d8abcc1cb1dda8d256acabcb9102795e42058f Mon Sep 17 00:00:00 2001 From: Serhii Skoryi <86241279+sskoryi-256@users.noreply.github.com> Date: Thu, 12 Sep 2024 15:40:48 +0300 Subject: [PATCH 12/12] Notification center: migrate from Office 365 Connectors to Microsoft Teams Workflows (#11583) * PROD-4064: Migrate from Office 365 Connectors to Microsoft Teams Workflows. Replace MessageCard by AdaptiveCard * PROD-4064: Migrate from Office 365 Connectors to Microsoft Teams Workflows. Replace MessageCard by AdaptiveCard * PROD-4064: Migrate from Office 365 Connectors to Microsoft Teams Workflows. Replace MessageCard by AdaptiveCard * PROD-4064: Migrate from Office 365 Connectors to Microsoft Teams Workflows. Replace MessageCard by AdaptiveCard * PROD-4064: Migrate from Office 365 Connectors to Microsoft Teams Workflows. Replace MessageCard by AdaptiveCard * PROD-4064: Migrate from Office 365 Connectors to Microsoft Teams Workflows. Replace MessageCard by AdaptiveCard * PROD-4064: Migrate from Office 365 Connectors to Microsoft Teams Workflows. Replace MessageCard by AdaptiveCard * UI: Change configuration Microsoft Team recipient config * UI: Add null check in Microsoft Team recipient config * PROD-4064: Migrate from Office 365 Connectors to Microsoft Teams Workflows. Replace MessageCard by AdaptiveCard * PROD-4064: Migrate from Office 365 Connectors to Microsoft Teams Workflows. Replace MessageCard by AdaptiveCard * PROD-4064: Migrate from Office 365 Connectors to Microsoft Teams Workflows. Replace MessageCard by AdaptiveCard * PROD-4064: Migrate from Office 365 Connectors to Microsoft Teams Workflows. Replace MessageCard by AdaptiveCard * PROD-4064: Migrate from Office 365 Connectors to Microsoft Teams Workflows. Replace MessageCard by AdaptiveCard * PROD-4064: Migrate from Office 365 Connectors to Microsoft Teams Workflows. Replace MessageCard by AdaptiveCard * PROD-4064: Migrate from Office 365 Connectors to Microsoft Teams Workflows. Replace MessageCard by AdaptiveCard * Resolved PR comments * Resolved PR comments --------- Co-authored-by: sskoryi Co-authored-by: Vladyslav_Prykhodko --- .../MicrosoftTeamsNotificationChannel.java | 179 +++++++++--------- .../channels/TeamsAdaptiveCard.java | 93 +++++++++ .../channels/TeamsMessageCard.java | 92 +++++++++ .../notification/NotificationApiTest.java | 86 ++++++++- ...icrosoftTeamsNotificationTargetConfig.java | 1 + .../server/dao/util/ImageUtils.java | 90 +++++++++ ...cipient-notification-dialog.component.html | 13 ++ ...recipient-notification-dialog.component.ts | 8 +- .../app/shared/models/notification.models.ts | 1 + .../assets/locale/locale.constant-en_US.json | 2 + 10 files changed, 470 insertions(+), 95 deletions(-) create mode 100644 application/src/main/java/org/thingsboard/server/service/notification/channels/TeamsAdaptiveCard.java create mode 100644 application/src/main/java/org/thingsboard/server/service/notification/channels/TeamsMessageCard.java diff --git a/application/src/main/java/org/thingsboard/server/service/notification/channels/MicrosoftTeamsNotificationChannel.java b/application/src/main/java/org/thingsboard/server/service/notification/channels/MicrosoftTeamsNotificationChannel.java index d3113787e5..358e2d1f6d 100644 --- a/application/src/main/java/org/thingsboard/server/service/notification/channels/MicrosoftTeamsNotificationChannel.java +++ b/application/src/main/java/org/thingsboard/server/service/notification/channels/MicrosoftTeamsNotificationChannel.java @@ -15,16 +15,17 @@ */ package org.thingsboard.server.service.notification.channels; -import com.fasterxml.jackson.annotation.JsonInclude; -import com.fasterxml.jackson.annotation.JsonProperty; +import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.databind.node.ObjectNode; import com.google.common.base.Strings; -import lombok.Data; import lombok.RequiredArgsConstructor; import lombok.Setter; import org.apache.commons.codec.binary.Base64; import org.apache.commons.lang3.StringUtils; import org.springframework.boot.web.client.RestTemplateBuilder; +import org.springframework.http.HttpEntity; +import org.springframework.http.HttpHeaders; +import org.springframework.http.MediaType; import org.springframework.stereotype.Component; import org.springframework.web.client.RestTemplate; import org.thingsboard.common.util.JacksonUtil; @@ -37,6 +38,8 @@ import org.thingsboard.server.common.data.notification.template.MicrosoftTeamsDe import org.thingsboard.server.service.notification.NotificationProcessingContext; import org.thingsboard.server.service.security.system.SystemSecurityService; +import java.net.URI; +import java.net.URISyntaxException; import java.time.Duration; import java.time.temporal.ChronoUnit; import java.util.List; @@ -56,17 +59,93 @@ public class MicrosoftTeamsNotificationChannel implements NotificationChannel request = new HttpEntity<>(JacksonUtil.toString(teamsAdaptiveCard), headers); + restTemplate.postForEntity(new URI(targetConfig.getWebhookUrl()), request, String.class); + } + + private void sendTeamsMessageCard(MicrosoftTeamsNotificationTargetConfig targetConfig, MicrosoftTeamsDeliveryMethodNotificationTemplate processedTemplate, NotificationProcessingContext ctx) throws JsonProcessingException, URISyntaxException { + TeamsMessageCard teamsMessageCard = new TeamsMessageCard(); + teamsMessageCard.setThemeColor(Strings.emptyToNull(processedTemplate.getThemeColor())); + if (StringUtils.isEmpty(processedTemplate.getSubject())) { + teamsMessageCard.setText(processedTemplate.getBody()); + } else { + teamsMessageCard.setSummary(processedTemplate.getSubject()); + TeamsMessageCard.Section section = new TeamsMessageCard.Section(); section.setActivityTitle(processedTemplate.getSubject()); section.setActivitySubtitle(processedTemplate.getBody()); - message.setSections(List.of(section)); + teamsMessageCard.setSections(List.of(section)); } + var button = processedTemplate.getButton(); + String uri = getButtonUri(processedTemplate, ctx); + + if (StringUtils.isNotBlank(uri) && button.getText() != null) { + TeamsMessageCard.ActionCard actionCard = new TeamsMessageCard.ActionCard(); + actionCard.setType("OpenUri"); + actionCard.setName(button.getText()); + var target = new TeamsMessageCard.ActionCard.Target("default", uri); + actionCard.setTargets(List.of(target)); + teamsMessageCard.setPotentialAction(List.of(actionCard)); + } + + HttpHeaders headers = new HttpHeaders(); + headers.setContentType(MediaType.APPLICATION_JSON); + HttpEntity request = new HttpEntity<>(JacksonUtil.toString(teamsMessageCard), headers); + restTemplate.postForEntity(new URI(targetConfig.getWebhookUrl()), request, String.class); + } + + private String getButtonUri(MicrosoftTeamsDeliveryMethodNotificationTemplate processedTemplate, NotificationProcessingContext ctx) throws JsonProcessingException { var button = processedTemplate.getButton(); if (button != null && button.isEnabled()) { String uri; @@ -99,17 +178,9 @@ public class MicrosoftTeamsNotificationChannel implements NotificationChannel sections; - private List potentialAction; - - @Data - public static class Section { - private String activityTitle; - private String activitySubtitle; - private String activityImage; - private List facts; - private boolean markdown; - - @Data - public static class Fact { - private final String name; - private final String value; - } - } - - @Data - @JsonInclude(JsonInclude.Include.NON_NULL) - public static class ActionCard { - @JsonProperty("@type") - private String type; // ActionCard, OpenUri - private String name; - private List inputs; // for ActionCard - private List actions; // for ActionCard - private List targets; - - @Data - public static class Input { - @JsonProperty("@type") - private String type; // TextInput, DateInput, MultichoiceInput - private String id; - private boolean isMultiple; - private String title; - private boolean isMultiSelect; - - @Data - public static class Choice { - private final String display; - private final String value; - } - } - - @Data - public static class Action { - @JsonProperty("@type") - private final String type; // HttpPOST - private final String name; - private final String target; // url - } - - @Data - public static class Target { - private final String os; - private final String uri; - } - } - - } - } diff --git a/application/src/main/java/org/thingsboard/server/service/notification/channels/TeamsAdaptiveCard.java b/application/src/main/java/org/thingsboard/server/service/notification/channels/TeamsAdaptiveCard.java new file mode 100644 index 0000000000..b5d434bac0 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/notification/channels/TeamsAdaptiveCard.java @@ -0,0 +1,93 @@ +/** + * Copyright © 2016-2024 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.service.notification.channels; + +import com.fasterxml.jackson.annotation.JsonProperty; +import lombok.AllArgsConstructor; +import lombok.Data; +import lombok.NoArgsConstructor; +import org.thingsboard.server.dao.util.ImageUtils; + +import java.util.ArrayList; +import java.util.List; + +/** + * @link AdaptiveCard Designer + */ +@Data +@NoArgsConstructor +@AllArgsConstructor +public class TeamsAdaptiveCard { + private String type = "message"; + private List attachments; + + @Data + @NoArgsConstructor + @AllArgsConstructor + public static class Attachment { + private String contentType = "application/vnd.microsoft.card.adaptive"; + private AdaptiveCard content; + } + + @Data + @NoArgsConstructor + @AllArgsConstructor + public static class AdaptiveCard { + @JsonProperty("$schema") + private final String schema = "http://adaptivecards.io/schemas/adaptive-card.json"; + private final String type = "AdaptiveCard"; + private BackgroundImage backgroundImage; + @JsonProperty("body") + private List textBlocks = new ArrayList<>(); + private List actions = new ArrayList<>(); + } + + @Data + @NoArgsConstructor + public static class BackgroundImage { + private String url; + private final String fillMode = "repeat"; + + public BackgroundImage(String color) { + // This is the only one way how to specify color the custom color for the card + url = ImageUtils.getEmbeddedBase64EncodedImg(color); + } + + } + + @Data + @NoArgsConstructor + @AllArgsConstructor + public static class TextBlock { + private final String type = "TextBlock"; + private String text; + private String weight = "Normal"; + private String size = "Medium"; + private String spacing = "None"; + private String color = "#FFFFFF"; + private final boolean wrap = true; + } + + @Data + @NoArgsConstructor + @AllArgsConstructor + public static class ActionOpenUrl { + private final String type = "Action.OpenUrl"; + private String title; + private String url; + } + +} \ No newline at end of file diff --git a/application/src/main/java/org/thingsboard/server/service/notification/channels/TeamsMessageCard.java b/application/src/main/java/org/thingsboard/server/service/notification/channels/TeamsMessageCard.java new file mode 100644 index 0000000000..65cdfec765 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/notification/channels/TeamsMessageCard.java @@ -0,0 +1,92 @@ +/** + * Copyright © 2016-2024 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.service.notification.channels; + +import com.fasterxml.jackson.annotation.JsonInclude; +import com.fasterxml.jackson.annotation.JsonProperty; +import lombok.Data; + +import java.util.List; + +@Data +public class TeamsMessageCard { + @JsonProperty("@type") + private final String type = "MessageCard"; + @JsonProperty("@context") + private final String context = "http://schema.org/extensions"; + private String themeColor; + private String summary; + private String text; + private List
sections; + private List potentialAction; + + @Data + public static class Section { + private String activityTitle; + private String activitySubtitle; + private String activityImage; + private List facts; + private boolean markdown; + + @Data + public static class Fact { + private final String name; + private final String value; + } + } + + @Data + @JsonInclude(JsonInclude.Include.NON_NULL) + public static class ActionCard { + @JsonProperty("@type") + private String type; // ActionCard, OpenUri + private String name; + private List inputs; // for ActionCard + private List actions; // for ActionCard + private List targets; + + @Data + public static class Input { + @JsonProperty("@type") + private String type; // TextInput, DateInput, MultichoiceInput + private String id; + private boolean isMultiple; + private String title; + private boolean isMultiSelect; + + @Data + public static class Choice { + private final String display; + private final String value; + } + } + + @Data + public static class Action { + @JsonProperty("@type") + private final String type; // HttpPOST + private final String name; + private final String target; // url + } + + @Data + public static class Target { + private final String os; + private final String uri; + } + } + +} diff --git a/application/src/test/java/org/thingsboard/server/service/notification/NotificationApiTest.java b/application/src/test/java/org/thingsboard/server/service/notification/NotificationApiTest.java index 53ae5f36ae..8bae1feca9 100644 --- a/application/src/test/java/org/thingsboard/server/service/notification/NotificationApiTest.java +++ b/application/src/test/java/org/thingsboard/server/service/notification/NotificationApiTest.java @@ -25,6 +25,7 @@ import org.junit.Test; import org.mockito.ArgumentCaptor; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.test.mock.mockito.MockBean; +import org.springframework.http.HttpEntity; import org.springframework.test.web.servlet.ResultActions; import org.springframework.web.client.RestTemplate; import org.thingsboard.common.util.JacksonUtil; @@ -85,8 +86,12 @@ import org.thingsboard.server.common.data.security.Authority; import org.thingsboard.server.dao.notification.DefaultNotifications; import org.thingsboard.server.dao.service.DaoSqlTest; import org.thingsboard.server.service.notification.channels.MicrosoftTeamsNotificationChannel; +import org.thingsboard.server.service.notification.channels.TeamsAdaptiveCard; +import org.thingsboard.server.service.notification.channels.TeamsMessageCard; import org.thingsboard.server.service.ws.notification.cmd.UnreadNotificationsUpdate; +import java.net.URI; +import java.net.URISyntaxException; import java.util.ArrayList; import java.util.HashMap; import java.util.List; @@ -752,7 +757,7 @@ public class NotificationApiTest extends AbstractNotificationApiTest { } @Test - public void testMicrosoftTeamsNotifications() throws Exception { + public void testMicrosoftTeamsNotificationsWithOfficeConnector() throws URISyntaxException { RestTemplate restTemplate = mock(RestTemplate.class); microsoftTeamsNotificationChannel.setRestTemplate(restTemplate); @@ -760,6 +765,7 @@ public class NotificationApiTest extends AbstractNotificationApiTest { var targetConfig = new MicrosoftTeamsNotificationTargetConfig(); targetConfig.setWebhookUrl(webhookUrl); targetConfig.setChannelName("My channel"); + targetConfig.setUseOldApi(true); NotificationTarget target = new NotificationTarget(); target.setName("Microsoft Teams channel"); target.setConfiguration(targetConfig); @@ -770,7 +776,7 @@ public class NotificationApiTest extends AbstractNotificationApiTest { String templateParams = "${recipientTitle} - ${entityType}"; template.setSubject("Subject: " + templateParams); template.setBody("Body: " + templateParams); - template.setThemeColor("ff0000"); + template.setThemeColor("#ff0000"); var button = new MicrosoftTeamsDeliveryMethodNotificationTemplate.Button(); button.setEnabled(true); button.setText("Button: " + templateParams); @@ -803,11 +809,13 @@ public class NotificationApiTest extends AbstractNotificationApiTest { assertThat(preview.getRecipientsCountByTarget().get(target.getName())).isEqualTo(1); assertThat(preview.getRecipientsPreview()).containsOnly(targetConfig.getChannelName()); - var messageCaptor = ArgumentCaptor.forClass(MicrosoftTeamsNotificationChannel.Message.class); + ArgumentCaptor> messageCaptor = ArgumentCaptor.forClass(HttpEntity.class); notificationCenter.processNotificationRequest(tenantId, notificationRequest, null); - verify(restTemplate, timeout(20000)).postForEntity(eq(webhookUrl), messageCaptor.capture(), any()); + verify(restTemplate, timeout(20000)).postForEntity(eq(new URI(webhookUrl)), messageCaptor.capture(), any()); + + HttpEntity value = messageCaptor.getValue(); + TeamsMessageCard message = JacksonUtil.fromString(value.getBody(), TeamsMessageCard.class); - var message = messageCaptor.getValue(); String expectedParams = "My channel - Device"; assertThat(message.getThemeColor()).isEqualTo(template.getThemeColor()); assertThat(message.getSections().get(0).getActivityTitle()).isEqualTo("Subject: " + expectedParams); @@ -816,6 +824,74 @@ public class NotificationApiTest extends AbstractNotificationApiTest { assertThat(message.getPotentialAction().get(0).getTargets().get(0).getUri()).isEqualTo("https://" + expectedParams); } + @Test + public void testMicrosoftTeamsNotificationsWithWorkflow() throws Exception { + RestTemplate restTemplate = mock(RestTemplate.class); + microsoftTeamsNotificationChannel.setRestTemplate(restTemplate); + + String webhookUrl = "https://webhook.com/webhookb2/9628fa60-d873-11ed-913c-a196b1f9b445"; + var targetConfig = new MicrosoftTeamsNotificationTargetConfig(); + targetConfig.setWebhookUrl(webhookUrl); + targetConfig.setChannelName("My channel"); + targetConfig.setUseOldApi(false); + NotificationTarget target = new NotificationTarget(); + target.setName("Microsoft Teams channel"); + target.setConfiguration(targetConfig); + target = saveNotificationTarget(target); + + var template = new MicrosoftTeamsDeliveryMethodNotificationTemplate(); + template.setEnabled(true); + String templateParams = "${recipientTitle} - ${entityType}"; + template.setSubject("Subject: " + templateParams); + template.setBody("Body: " + templateParams); + template.setThemeColor("#ff0000"); + var button = new MicrosoftTeamsDeliveryMethodNotificationTemplate.Button(); + button.setEnabled(true); + button.setText("Button: " + templateParams); + button.setLinkType(LinkType.LINK); + button.setLink("https://" + templateParams); + template.setButton(button); + NotificationTemplate notificationTemplate = new NotificationTemplate(); + notificationTemplate.setName("Notification to Teams"); + notificationTemplate.setNotificationType(NotificationType.GENERAL); + NotificationTemplateConfig templateConfig = new NotificationTemplateConfig(); + templateConfig.setDeliveryMethodsTemplates(Map.of( + NotificationDeliveryMethod.MICROSOFT_TEAMS, template + )); + notificationTemplate.setConfiguration(templateConfig); + notificationTemplate = saveNotificationTemplate(notificationTemplate); + + NotificationRequest notificationRequest = NotificationRequest.builder() + .tenantId(tenantId) + .originatorEntityId(tenantAdminUserId) + .templateId(notificationTemplate.getId()) + .targets(List.of(target.getUuidId())) + .info(EntityActionNotificationInfo.builder() + .entityId(new DeviceId(UUID.randomUUID())) + .actionType(ActionType.ADDED) + .userId(tenantAdminUserId.getId()) + .build()) + .build(); + + NotificationRequestPreview preview = doPost("/api/notification/request/preview", notificationRequest, NotificationRequestPreview.class); + assertThat(preview.getRecipientsCountByTarget().get(target.getName())).isEqualTo(1); + assertThat(preview.getRecipientsPreview()).containsOnly(targetConfig.getChannelName()); + + ArgumentCaptor> messageCaptor = ArgumentCaptor.forClass(HttpEntity.class); + notificationCenter.processNotificationRequest(tenantId, notificationRequest, null); + verify(restTemplate, timeout(20000)).postForEntity(eq(new URI(webhookUrl)), messageCaptor.capture(), any()); + + HttpEntity value = messageCaptor.getValue(); + TeamsAdaptiveCard message = JacksonUtil.fromString(value.getBody(), TeamsAdaptiveCard.class); + String expectedParams = "My channel - Device"; + assertThat(message).isNotNull(); + assertThat(message.getAttachments().get(0).getContent().getBackgroundImage().getUrl()).isNotEmpty(); + assertThat(message.getAttachments().get(0).getContent().getTextBlocks().get(0).getText()).isEqualTo("Subject: " + expectedParams); + assertThat(message.getAttachments().get(0).getContent().getTextBlocks().get(1).getText()).isEqualTo("Body: " + expectedParams); + assertThat(message.getAttachments().get(0).getContent().getActions().get(0).getTitle()).isEqualTo("Button: " + expectedParams); + assertThat(message.getAttachments().get(0).getContent().getActions().get(0).getUrl()).isEqualTo("https://" + expectedParams); + } + @Test public void testMobileAppNotifications() throws Exception { loginCustomerUser(); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/notification/targets/MicrosoftTeamsNotificationTargetConfig.java b/common/data/src/main/java/org/thingsboard/server/common/data/notification/targets/MicrosoftTeamsNotificationTargetConfig.java index e07abe7c00..dbd09698f0 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/notification/targets/MicrosoftTeamsNotificationTargetConfig.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/notification/targets/MicrosoftTeamsNotificationTargetConfig.java @@ -28,6 +28,7 @@ public class MicrosoftTeamsNotificationTargetConfig extends NotificationTargetCo private String webhookUrl; @NotEmpty private String channelName; + private Boolean useOldApi = Boolean.TRUE; @Override public NotificationTargetType getType() { diff --git a/dao/src/main/java/org/thingsboard/server/dao/util/ImageUtils.java b/dao/src/main/java/org/thingsboard/server/dao/util/ImageUtils.java index e92adba01d..19649655a3 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/util/ImageUtils.java +++ b/dao/src/main/java/org/thingsboard/server/dao/util/ImageUtils.java @@ -46,9 +46,11 @@ import org.thingsboard.server.common.data.StringUtils; import org.w3c.dom.Document; import javax.imageio.ImageIO; +import java.awt.Color; import java.awt.image.BufferedImage; import java.io.ByteArrayInputStream; import java.io.ByteArrayOutputStream; +import java.util.Base64; import java.util.Map; @NoArgsConstructor(access = AccessLevel.PRIVATE) @@ -265,6 +267,94 @@ public class ImageUtils { return new int[]{thumbnailWidth, thumbnailHeight}; } + public static String getEmbeddedBase64EncodedImg(String colorStr) { + try { + Color color = parseColor(colorStr); // Support for hex, rgb, hsla + BufferedImage image = new BufferedImage(1, 1, BufferedImage.TYPE_INT_RGB); + image.setRGB(0, 0, color.getRGB()); + + ByteArrayOutputStream outputStream = new ByteArrayOutputStream(); + ImageIO.write(image, "png", outputStream); + byte[] imageBytes = outputStream.toByteArray(); + String base64String = Base64.getEncoder().encodeToString(imageBytes); + + return "data:image/png;base64," + base64String; + } catch (Exception e) { + log.warn("Failed to generate embedded image for color: {}", colorStr, e); + return null; + } + } + + private static Color parseColor(String colorStr) { + if (colorStr.startsWith("#")) { + return Color.decode(colorStr); + } + + if (colorStr.startsWith("rgb")) { + return parseRgbColor(colorStr); + } + + if (colorStr.startsWith("hsl")) { + return parseHslaColor(colorStr); + } + + throw new IllegalArgumentException("Unsupported color format: " + colorStr); + } + + private static Color parseRgbColor(String rgb) { + String[] rgbValues = rgb.replaceAll("[^0-9,]", "").split(","); + int r = Integer.parseInt(rgbValues[0]); + int g = Integer.parseInt(rgbValues[1]); + int b = Integer.parseInt(rgbValues[2]); + return new Color(r, g, b); + } + + private static Color parseHslaColor(String hsla) { + String[] hslaValues = hsla.replaceAll("[^0-9.,]", "").split(","); + float h = Float.parseFloat(hslaValues[0]); + float s = Float.parseFloat(hslaValues[1]) / 100; + float l = Float.parseFloat(hslaValues[2]) / 100; + float a = hslaValues.length > 3 ? Float.parseFloat(hslaValues[3]) : 1.0f; + return hslaToColor(h, s, l, a); + } + + private static Color hslaToColor(float h, float s, float l, float alpha) { + float c = (1 - Math.abs(2 * l - 1)) * s; + float x = c * (1 - Math.abs((h / 60) % 2 - 1)); + float m = l - c / 2; + + float r = 0, g = 0, b = 0; + if (h < 60) { + r = c; + g = x; + } else if (h < 120) { + r = x; + g = c; + } else if (h < 180) { + g = c; + b = x; + } else if (h < 240) { + g = x; + b = c; + } else if (h < 300) { + r = x; + b = c; + } else { + r = c; + b = x; + } + + r += m; + g += m; + b += m; + + return new Color(clamp(r), clamp(g), clamp(b), clamp(alpha)); + } + + private static float clamp(float value) { + return Math.max(0, Math.min(1, value)); + } + @Data @AllArgsConstructor @NoArgsConstructor diff --git a/ui-ngx/src/app/modules/home/pages/notification/recipient/recipient-notification-dialog.component.html b/ui-ngx/src/app/modules/home/pages/notification/recipient/recipient-notification-dialog.component.html index 40cfe51e98..311dbee6d9 100644 --- a/ui-ngx/src/app/modules/home/pages/notification/recipient/recipient-notification-dialog.component.html +++ b/ui-ngx/src/app/modules/home/pages/notification/recipient/recipient-notification-dialog.component.html @@ -121,6 +121,19 @@
+
+ + {{ "notification.use-old-api" | translate }} + + + open_in_new + +
notification.webhook-url diff --git a/ui-ngx/src/app/modules/home/pages/notification/recipient/recipient-notification-dialog.component.ts b/ui-ngx/src/app/modules/home/pages/notification/recipient/recipient-notification-dialog.component.ts index 0a4be127ad..ad7c50d604 100644 --- a/ui-ngx/src/app/modules/home/pages/notification/recipient/recipient-notification-dialog.component.ts +++ b/ui-ngx/src/app/modules/home/pages/notification/recipient/recipient-notification-dialog.component.ts @@ -32,7 +32,7 @@ import { MAT_DIALOG_DATA, MatDialogRef } from '@angular/material/dialog'; import { FormBuilder, FormGroup, Validators } from '@angular/forms'; import { NotificationService } from '@core/http/notification.service'; import { EntityType } from '@shared/models/entity-type.models'; -import { deepTrim, isDefinedAndNotNull } from '@core/utils'; +import { deepTrim, isDefinedAndNotNull, isUndefinedOrNull } from '@core/utils'; import { Subject } from 'rxjs'; import { takeUntil } from 'rxjs/operators'; import { Authority } from '@shared/models/authority.enum'; @@ -100,6 +100,7 @@ export class RecipientNotificationDialogComponent extends conversation: [{value: '', disabled: true}, Validators.required], webhookUrl: [{value: '', disabled: true}, Validators.required], channelName: [{value: '', disabled: true}, Validators.required], + useOldApi: [{value: !this.isAdd, disabled: true}], description: [null] }) }); @@ -120,6 +121,7 @@ export class RecipientNotificationDialogComponent extends case NotificationTargetType.MICROSOFT_TEAMS: this.targetNotificationForm.get('configuration.webhookUrl').enable({emitEvent: false}); this.targetNotificationForm.get('configuration.channelName').enable({emitEvent: false}); + this.targetNotificationForm.get('configuration.useOldApi').enable({emitEvent: false}); break; } this.targetNotificationForm.get('configuration.type').enable({emitEvent: false}); @@ -169,6 +171,10 @@ export class RecipientNotificationDialogComponent extends this.targetNotificationForm.get('configuration.usersFilter.filterByTenants') .patchValue(!Array.isArray(this.data.target.configuration.usersFilter.tenantProfilesIds), {onlySelf: true}); } + if (data.target.configuration.type === NotificationTargetType.MICROSOFT_TEAMS + && isUndefinedOrNull(this.data.target.configuration.useOldApi)) { + this.targetNotificationForm.get('configuration.useOldApi').patchValue(true, {emitEvent: false}); + } } } diff --git a/ui-ngx/src/app/shared/models/notification.models.ts b/ui-ngx/src/app/shared/models/notification.models.ts index cfb1312209..dd5550e907 100644 --- a/ui-ngx/src/app/shared/models/notification.models.ts +++ b/ui-ngx/src/app/shared/models/notification.models.ts @@ -290,6 +290,7 @@ export interface SlackNotificationTargetConfig { export interface MicrosoftTeamsNotificationTargetConfig { webhookUrl: string; channelName: string; + useOldApi?: boolean; } export enum NotificationTargetType { PLATFORM_USERS = 'PLATFORM_USERS', diff --git a/ui-ngx/src/assets/locale/locale.constant-en_US.json b/ui-ngx/src/assets/locale/locale.constant-en_US.json index 37ea4a1a87..1a2089823c 100644 --- a/ui-ngx/src/assets/locale/locale.constant-en_US.json +++ b/ui-ngx/src/assets/locale/locale.constant-en_US.json @@ -4295,6 +4295,8 @@ "type": "Type", "unread": "Unread", "updated": "Updated", + "use-deprecated-webhook-connectors": "Use deprecated Webhook connectors", + "use-old-api": "Use old API", "use-template": "Use template", "view-all": "View all", "warning": "Warning",