Browse Source

add enqueueForTellNext

pull/7244/head
Yuriy Lytvynchuk 4 years ago
parent
commit
d41a791dc3
  1. 3
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbAbstractTransformNode.java
  2. 20
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbSplitArrayMsgNode.java
  3. 36
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/transform/TbSplitArrayMsgNodeTest.java

3
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);

20
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!"));

36
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<TbMsg> 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<Runnable> successCaptor = ArgumentCaptor.forClass(Runnable.class);
ArgumentCaptor<Consumer<Throwable>> 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<TbMsg> newMsgCaptor = ArgumentCaptor.forClass(TbMsg.class);
verify(ctx, times(dataNode.size())).tellSuccess(newMsgCaptor.capture());
}
verify(ctx, never()).tellFailure(any(), any());
}

Loading…
Cancel
Save