From d41a791dc312a0cb543df986ca13d3f8fd60171d Mon Sep 17 00:00:00 2001 From: Yuriy Lytvynchuk Date: Tue, 27 Sep 2022 17:12:01 +0300 Subject: [PATCH] add enqueueForTellNext --- .../transform/TbAbstractTransformNode.java | 3 +- .../engine/transform/TbSplitArrayMsgNode.java | 20 ++++++++--- .../transform/TbSplitArrayMsgNodeTest.java | 36 +++++++++++-------- 3 files changed, 39 insertions(+), 20 deletions(-) diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbAbstractTransformNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbAbstractTransformNode.java index 29c61b81c7..e51a2d9205 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbAbstractTransformNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbAbstractTransformNode.java @@ -22,6 +22,7 @@ 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.rule.engine.api.TbRelationTypes; import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.queue.RuleEngineException; @@ -81,7 +82,7 @@ public abstract class TbAbstractTransformNode implements TbNode { ctx.tellFailure(msg, e); } }); - msgs.forEach(newMsg -> ctx.enqueueForTellNext(newMsg, "Success", wrapper::onSuccess, wrapper::onFailure)); + msgs.forEach(newMsg -> ctx.enqueueForTellNext(newMsg, TbRelationTypes.SUCCESS, wrapper::onSuccess, wrapper::onFailure)); } } else { ctx.tellNext(msg, FAILURE); diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbSplitArrayMsgNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbSplitArrayMsgNode.java index 965cc3f058..2dc27b6f92 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbSplitArrayMsgNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbSplitArrayMsgNode.java @@ -25,9 +25,12 @@ 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.rule.engine.api.TbRelationTypes; import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.msg.TbMsg; +import org.thingsboard.server.common.msg.queue.RuleEngineException; +import org.thingsboard.server.common.msg.queue.TbMsgCallback; import java.util.concurrent.ExecutionException; @@ -46,7 +49,7 @@ import java.util.concurrent.ExecutionException; ) public class TbSplitArrayMsgNode implements TbNode { - EmptyNodeConfiguration config; + private EmptyNodeConfiguration config; @Override public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException { @@ -61,10 +64,19 @@ public class TbSplitArrayMsgNode implements TbNode { if (data.size() == 1) { ctx.tellSuccess(TbMsg.transformMsg(msg, msg.getType(), msg.getOriginator(), msg.getMetaData(), JacksonUtil.toString(data.get(0)))); } else { - ctx.ack(msg); - data.forEach(msgNode -> { - ctx.tellSuccess(TbMsg.newMsg(msg.getQueueName(), msg.getType(), msg.getOriginator(), msg.getMetaData(), JacksonUtil.toString(msgNode))); + TbMsgCallbackWrapper wrapper = new MultipleTbMsgsCallbackWrapper(data.size(), new TbMsgCallback() { + @Override + public void onSuccess() { + ctx.ack(msg); + } + + @Override + public void onFailure(RuleEngineException e) { + ctx.tellFailure(msg, e); + } }); + data.forEach(msgNode -> ctx.enqueueForTellNext(TbMsg.newMsg(msg.getQueueName(), msg.getType(), msg.getOriginator(), msg.getMetaData(), JacksonUtil.toString(msgNode)), + TbRelationTypes.SUCCESS, wrapper::onSuccess, wrapper::onFailure)); } } else { ctx.tellFailure(msg, new RuntimeException("Msg data is not a JSON Array!")); diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/transform/TbSplitArrayMsgNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/transform/TbSplitArrayMsgNodeTest.java index adaf842d42..2925ab96ae 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/transform/TbSplitArrayMsgNodeTest.java +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/transform/TbSplitArrayMsgNodeTest.java @@ -17,7 +17,6 @@ package org.thingsboard.rule.engine.transform; import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; -import com.fasterxml.jackson.databind.node.ArrayNode; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; @@ -35,9 +34,11 @@ import org.thingsboard.server.common.msg.queue.TbMsgCallback; import java.util.Map; import java.util.UUID; +import java.util.function.Consumer; import static org.assertj.core.api.Assertions.assertThat; import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.never; import static org.mockito.Mockito.spy; @@ -70,27 +71,22 @@ public class TbSplitArrayMsgNodeTest { node.destroy(); } - @Test - void givenDefaultConfig_whenInit_thenOK() { - assertThat(node.config).isEqualTo(config); - } - @Test void givenFewMsg_whenOnMsg_thenVerifyOutput() throws Exception { String data = "[{\"Attribute_1\":22.5,\"Attribute_2\":10.3}, {\"Attribute_1\":1,\"Attribute_2\":2}]"; - VerifyOutputMsg(data, 2); + VerifyOutputMsg(data); } @Test void givenOneMsg_whenOnMsg_thenVerifyOutput() throws Exception { String data = "[{\"Attribute_1\":22.5,\"Attribute_2\":10.3}]"; - VerifyOutputMsg(data, 1); + VerifyOutputMsg(data); } @Test void givenZeroMsg_whenOnMsg_thenVerifyOutput() throws Exception { String data = "[]"; - VerifyOutputMsg(data, 0); + VerifyOutputMsg(data); } @Test @@ -113,13 +109,23 @@ public class TbSplitArrayMsgNodeTest { assertThat(newMsg).isSameAs(msg); } - private void VerifyOutputMsg(String data, int sizeArray) throws Exception { + private void VerifyOutputMsg(String data) throws Exception { JsonNode dataNode = JacksonUtil.toJsonNode(data); - node.onMsg(ctx, getTbMsg(deviceId, dataNode.toString())); - - ArgumentCaptor newMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); - verify(ctx, times(dataNode.size())).tellSuccess(newMsgCaptor.capture()); - verify(ctx, times(sizeArray)).tellSuccess(newMsgCaptor.capture()); + TbMsg tbMsg = getTbMsg(deviceId, dataNode.toString()); + node.onMsg(ctx, tbMsg); + + if (dataNode.size() > 1) { + ArgumentCaptor successCaptor = ArgumentCaptor.forClass(Runnable.class); + ArgumentCaptor> failureCaptor = ArgumentCaptor.forClass(Consumer.class); + verify(ctx, times(dataNode.size())).enqueueForTellNext(any(), anyString(), successCaptor.capture(), failureCaptor.capture()); + for (Runnable valueCaptor : successCaptor.getAllValues()) { + valueCaptor.run(); + } + verify(ctx, times(1)).ack(tbMsg); + } else { + ArgumentCaptor newMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); + verify(ctx, times(dataNode.size())).tellSuccess(newMsgCaptor.capture()); + } verify(ctx, never()).tellFailure(any(), any()); }