Browse Source

Refactor TbMsg transforming

pull/12256/head
ViacheslavKlimov 2 years ago
parent
commit
0b3ffb5b4a
  1. 23
      application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java
  2. 11
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java
  3. 6
      application/src/main/java/org/thingsboard/server/service/script/RuleNodeJsScriptEngine.java
  4. 6
      application/src/main/java/org/thingsboard/server/service/script/RuleNodeTbelScriptEngine.java
  5. 2
      application/src/test/java/org/thingsboard/server/service/queue/DefaultTbClusterServiceTest.java
  6. 215
      common/message/src/main/java/org/thingsboard/server/common/msg/TbMsg.java
  7. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbAbstractAlarmNode.java
  8. 9
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/lambda/TbAwsLambdaNode.java
  9. 8
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/sns/TbSnsNode.java
  10. 8
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/sqs/TbSqsNode.java
  11. 8
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/gcp/pubsub/TbPubSubNode.java
  12. 8
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/kafka/TbKafkaNode.java
  13. 8
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/math/TbMathNode.java
  14. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/CalculateDeltaNode.java
  15. 8
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractNodeWithFetchTo.java
  16. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNode.java
  17. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mqtt/TbMqttNode.java
  18. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbNotificationNode.java
  19. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rabbitmq/TbRabbitMqNode.java
  20. 8
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbHttpClient.java
  21. 5
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbCopyKeysNode.java
  22. 5
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbDeleteKeysNode.java
  23. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbJsonPathNode.java
  24. 5
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbRenameKeysNode.java
  25. 8
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbSplitArrayMsgNode.java
  26. 72
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/action/TbCreateAlarmNodeTest.java
  27. 4
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/action/TbCreateRelationNodeTest.java
  28. 18
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/aws/lambda/TbAwsLambdaNodeTest.java
  29. 16
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/gcp/pubsub/TbPubSubNodeTest.java
  30. 8
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/kafka/TbKafkaNodeTest.java
  31. 8
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNodeTest.java
  32. 4
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/mqtt/TbMqttNodeTest.java
  33. 4
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/rabbitmq/TbRabbitMqNodeTest.java
  34. 16
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNodeTest.java

23
application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java

@ -377,7 +377,12 @@ public class DefaultTbContext implements TbContext {
@Override
public TbMsg transformMsg(TbMsg origMsg, String type, EntityId originator, TbMsgMetaData metaData, String data) {
return TbMsg.transformMsg(origMsg, type, originator, metaData, data);
return origMsg.transform()
.type(type)
.originator(originator)
.metaData(metaData)
.data(data)
.build();
}
@Override
@ -401,17 +406,27 @@ public class DefaultTbContext implements TbContext {
@Override
public TbMsg transformMsg(TbMsg origMsg, TbMsgType type, EntityId originator, TbMsgMetaData metaData, String data) {
return TbMsg.transformMsg(origMsg, type, originator, metaData, data);
return origMsg.transform()
.type(type)
.originator(originator)
.metaData(metaData)
.data(data)
.build();
}
@Override
public TbMsg transformMsg(TbMsg origMsg, TbMsgMetaData metaData, String data) {
return TbMsg.transformMsg(origMsg, metaData, data);
return origMsg.transform()
.metaData(metaData)
.data(data)
.build();
}
@Override
public TbMsg transformMsgOriginator(TbMsg origMsg, EntityId originator) {
return TbMsg.transformMsgOriginator(origMsg, originator);
return origMsg.transform()
.originator(originator)
.build();
}
@Override

11
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java

@ -290,11 +290,16 @@ public class DefaultTbClusterService implements TbClusterService {
boolean isQueueTransform = targetQueueName != null && !targetQueueName.equals(tbMsg.getQueueName());
if (isRuleChainTransform && isQueueTransform) {
tbMsg = TbMsg.transformMsg(tbMsg, targetRuleChainId, targetQueueName);
tbMsg = tbMsg.transform()
.queueName(targetQueueName)
.ruleChainId(targetRuleChainId)
.build();
} else if (isRuleChainTransform) {
tbMsg = TbMsg.transformMsgRuleChainId(tbMsg, targetRuleChainId);
tbMsg = tbMsg.transform()
.ruleChainId(targetRuleChainId)
.build();
} else if (isQueueTransform) {
tbMsg = TbMsg.transformMsgQueueName(tbMsg, targetQueueName);
tbMsg = tbMsg.transform(targetQueueName);
}
}
return tbMsg;

6
application/src/main/java/org/thingsboard/server/service/script/RuleNodeJsScriptEngine.java

@ -147,6 +147,10 @@ public class RuleNodeJsScriptEngine extends RuleNodeScriptEngine<JsInvokeService
String newData = data != null ? data : msg.getData();
TbMsgMetaData newMetadata = metadata != null ? new TbMsgMetaData(metadata) : msg.getMetaData().copy();
String newMessageType = !StringUtils.isEmpty(messageType) ? messageType : msg.getType();
return TbMsg.transformMsg(msg, newMessageType, msg.getOriginator(), newMetadata, newData);
return msg.transform()
.type(newMessageType)
.metaData(newMetadata)
.data(newData)
.build();
}
}

6
application/src/main/java/org/thingsboard/server/service/script/RuleNodeTbelScriptEngine.java

@ -156,7 +156,11 @@ public class RuleNodeTbelScriptEngine extends RuleNodeScriptEngine<TbelInvokeSer
String newData = data != null ? data : msg.getData();
TbMsgMetaData newMetadata = metadata != null ? new TbMsgMetaData(metadata) : msg.getMetaData().copy();
String newMessageType = !StringUtils.isEmpty(messageType) ? messageType : msg.getType();
return TbMsg.transformMsg(msg, newMessageType, msg.getOriginator(), newMetadata, newData);
return msg.transform()
.type(newMessageType)
.metaData(newMetadata)
.data(newData)
.build();
}
private static <T> ListenableFuture<T> wrongResultType(Object result) {

2
application/src/test/java/org/thingsboard/server/service/queue/DefaultTbClusterServiceTest.java

@ -371,7 +371,7 @@ public class DefaultTbClusterServiceTest {
clusterService.pushMsgToRuleEngine(tenantId, deviceId, requestMsg, false, callback);
verify(producerProvider).getRuleEngineMsgProducer();
TbMsg expectedMsg = TbMsg.transformMsgQueueName(requestMsg, DataConstants.MAIN_QUEUE_NAME);
TbMsg expectedMsg = requestMsg.transform(DataConstants.MAIN_QUEUE_NAME);
ArgumentCaptor<TbMsg> actualMsg = ArgumentCaptor.forClass(TbMsg.class);
verify(ruleEngineProducerService).sendToRuleEngine(eq(tbREQueueProducer), eq(tenantId), actualMsg.capture(), eq(callback));
assertThat(actualMsg.getValue()).usingRecursiveComparison().ignoringFields("ctx").isEqualTo(expectedMsg);

215
common/message/src/main/java/org/thingsboard/server/common/msg/TbMsg.java

@ -81,82 +81,42 @@ public final class TbMsg implements Serializable {
return ctx.getAndIncrementRuleNodeCounter();
}
/**
* Transforms an existing TbMsg instance by changing its message type, originator, metadata, and data.
*
* <p><strong>Deprecated:</strong> This method is deprecated since version 3.6.0 and should only be used when you need to
* specify a custom message type that doesn't exist in the {@link TbMsgType} enum. For standard message types,
* it is recommended to use the {@link #transformMsg(TbMsg, TbMsgType, EntityId, TbMsgMetaData, String)}
* method instead.</p>
*
*
* @param tbMsg the TbMsg instance to transform
* @param type the new message type
* @param originator the new originator
* @param metaData the new metadata
* @param data the new data
* @return the transformed TbMsg instance
*/
@Deprecated(since = "3.6.0")
public static TbMsg transformMsg(TbMsg tbMsg, String type, EntityId originator, TbMsgMetaData metaData, String data) {
return tbMsg.transform()
.type(type)
.originator(originator)
.metaData(metaData)
.data(data)
.ctx(tbMsg.ctx)
public TbMsg transform(String queueName) {
return transform()
.queueName(queueName)
.ruleNodeId(null)
.build();
}
public static TbMsg transformMsg(TbMsg tbMsg, TbMsgType type, EntityId originator, TbMsgMetaData metaData, String data) {
return new TbMsg(tbMsg.queueName, tbMsg.id, tbMsg.ts, type, type.name(), originator, tbMsg.customerId, metaData.copy(), tbMsg.dataType,
data, tbMsg.ruleChainId, tbMsg.ruleNodeId, tbMsg.correlationId, tbMsg.partition, tbMsg.ctx.copy(), tbMsg.callback);
}
public static TbMsg transformMsgOriginator(TbMsg tbMsg, EntityId originatorId) {
return new TbMsg(tbMsg.queueName, tbMsg.id, tbMsg.ts, tbMsg.internalType, tbMsg.type, originatorId, tbMsg.getCustomerId(), tbMsg.metaData, tbMsg.dataType,
tbMsg.data, tbMsg.ruleChainId, tbMsg.ruleNodeId, tbMsg.correlationId, tbMsg.partition, tbMsg.ctx.copy(), tbMsg.getCallback());
}
public static TbMsg transformMsgData(TbMsg tbMsg, String data) {
return new TbMsg(tbMsg.queueName, tbMsg.id, tbMsg.ts, tbMsg.internalType, tbMsg.type, tbMsg.originator, tbMsg.customerId, tbMsg.metaData, tbMsg.dataType,
data, tbMsg.ruleChainId, tbMsg.ruleNodeId, tbMsg.correlationId, tbMsg.partition, tbMsg.ctx.copy(), tbMsg.getCallback());
}
public static TbMsg transformMsgMetadata(TbMsg tbMsg, TbMsgMetaData metadata) {
return new TbMsg(tbMsg.queueName, tbMsg.id, tbMsg.ts, tbMsg.internalType, tbMsg.type, tbMsg.originator, tbMsg.customerId, metadata.copy(), tbMsg.dataType,
tbMsg.data, tbMsg.ruleChainId, tbMsg.ruleNodeId, tbMsg.correlationId, tbMsg.partition, tbMsg.ctx.copy(), tbMsg.getCallback());
}
public static TbMsg transformMsg(TbMsg tbMsg, TbMsgMetaData metadata, String data) {
return new TbMsg(tbMsg.queueName, tbMsg.id, tbMsg.ts, tbMsg.internalType, tbMsg.type, tbMsg.originator, tbMsg.customerId, metadata, tbMsg.dataType,
data, tbMsg.ruleChainId, tbMsg.ruleNodeId, tbMsg.correlationId, tbMsg.partition, tbMsg.ctx.copy(), tbMsg.getCallback());
}
public static TbMsg transformMsgCustomerId(TbMsg tbMsg, CustomerId customerId) {
return new TbMsg(tbMsg.queueName, tbMsg.id, tbMsg.ts, tbMsg.internalType, tbMsg.type, tbMsg.originator, customerId, tbMsg.metaData, tbMsg.dataType,
tbMsg.data, tbMsg.ruleChainId, tbMsg.ruleNodeId, tbMsg.correlationId, tbMsg.partition, tbMsg.ctx.copy(), tbMsg.getCallback());
// used for enqueueForTellNext
public static TbMsg newMsg(TbMsg tbMsg, String queueName, RuleChainId ruleChainId, RuleNodeId ruleNodeId) {
return tbMsg.transform()
.id(UUID.randomUUID())
.queueName(queueName)
.metaData(tbMsg.getMetaData())
.ruleChainId(ruleChainId)
.ruleNodeId(ruleNodeId)
.callback(TbMsgCallback.EMPTY)
.build();
}
public static TbMsg transformMsgRuleChainId(TbMsg tbMsg, RuleChainId ruleChainId) {
return new TbMsg(tbMsg.queueName, tbMsg.id, tbMsg.ts, tbMsg.internalType, tbMsg.type, tbMsg.originator, tbMsg.customerId, tbMsg.metaData, tbMsg.dataType,
tbMsg.data, ruleChainId, null, tbMsg.correlationId, tbMsg.partition, tbMsg.ctx.copy(), tbMsg.getCallback());
public TbMsg copyWithRuleChainId(RuleChainId ruleChainId) {
return copyWithRuleChainId(ruleChainId, this.id);
}
public static TbMsg transformMsgQueueName(TbMsg tbMsg, String queueName) {
return new TbMsg(queueName, tbMsg.id, tbMsg.ts, tbMsg.internalType, tbMsg.type, tbMsg.originator, tbMsg.customerId, tbMsg.metaData, tbMsg.dataType,
tbMsg.data, tbMsg.getRuleChainId(), null, tbMsg.correlationId, tbMsg.partition, tbMsg.ctx.copy(), tbMsg.getCallback());
public TbMsg copyWithRuleChainId(RuleChainId ruleChainId, UUID msgId) {
return new TbMsg(this.queueName, msgId, this.ts, this.internalType, this.type, this.originator, this.customerId,
this.metaData, this.dataType, this.data, ruleChainId, null, this.correlationId, this.partition, this.ctx, callback);
}
public static TbMsg transformMsg(TbMsg tbMsg, RuleChainId ruleChainId, String queueName) {
return new TbMsg(queueName, tbMsg.id, tbMsg.ts, tbMsg.internalType, tbMsg.type, tbMsg.originator, tbMsg.customerId, tbMsg.metaData, tbMsg.dataType,
tbMsg.data, ruleChainId, null, tbMsg.correlationId, tbMsg.partition, tbMsg.ctx.copy(), tbMsg.getCallback());
public TbMsg copyWithRuleNodeId(RuleChainId ruleChainId, RuleNodeId ruleNodeId, UUID msgId) {
return new TbMsg(this.queueName, msgId, this.ts, this.internalType, this.type, this.originator, this.customerId,
this.metaData, this.dataType, this.data, ruleChainId, ruleNodeId, this.correlationId, this.partition, this.ctx, callback);
}
//used for enqueueForTellNext
public static TbMsg newMsg(TbMsg tbMsg, String queueName, RuleChainId ruleChainId, RuleNodeId ruleNodeId) {
return new TbMsg(queueName, UUID.randomUUID(), tbMsg.getTs(), tbMsg.getInternalType(), tbMsg.getType(), tbMsg.getOriginator(), tbMsg.customerId, tbMsg.getMetaData().copy(),
tbMsg.getDataType(), tbMsg.getData(), ruleChainId, ruleNodeId, tbMsg.correlationId, tbMsg.partition, tbMsg.ctx.copy(), TbMsgCallback.EMPTY);
public TbMsg copyWithNewCtx() {
return new TbMsg(this.queueName, this.id, this.ts, this.internalType, this.type, this.originator, this.customerId,
this.metaData, this.dataType, this.data, ruleChainId, ruleNodeId, this.correlationId, this.partition, this.ctx.copy(), TbMsgCallback.EMPTY);
}
private TbMsg(String queueName, UUID id, long ts, TbMsgType internalType, String type, EntityId originator, CustomerId customerId, TbMsgMetaData metaData, TbMsgDataType dataType, String data,
@ -278,25 +238,6 @@ public final class TbMsg implements Serializable {
}
}
public TbMsg copyWithRuleChainId(RuleChainId ruleChainId) {
return copyWithRuleChainId(ruleChainId, this.id);
}
public TbMsg copyWithRuleChainId(RuleChainId ruleChainId, UUID msgId) {
return new TbMsg(this.queueName, msgId, this.ts, this.internalType, this.type, this.originator, this.customerId,
this.metaData, this.dataType, this.data, ruleChainId, null, this.correlationId, this.partition, this.ctx, callback);
}
public TbMsg copyWithRuleNodeId(RuleChainId ruleChainId, RuleNodeId ruleNodeId, UUID msgId) {
return new TbMsg(this.queueName, msgId, this.ts, this.internalType, this.type, this.originator, this.customerId,
this.metaData, this.dataType, this.data, ruleChainId, ruleNodeId, this.correlationId, this.partition, this.ctx, callback);
}
public TbMsg copyWithNewCtx() {
return new TbMsg(this.queueName, this.id, this.ts, this.internalType, this.type, this.originator, this.customerId,
this.metaData, this.dataType, this.data, ruleChainId, ruleNodeId, this.correlationId, this.partition, this.ctx.copy(), TbMsgCallback.EMPTY);
}
public TbMsgCallback getCallback() {
// May be null in case of deserialization;
return Objects.requireNonNullElse(callback, TbMsgCallback.EMPTY);
@ -357,62 +298,91 @@ public final class TbMsg implements Serializable {
}
public TbMsgBuilder transform() {
return new TbMsgTransformer()
.queueName(this.queueName)
.id(this.id)
.ts(this.ts)
.type(this.type)
.type(this.internalType)
.originator(this.originator)
.customerId(this.customerId)
.metaData(this.metaData)
.dataType(this.dataType)
.data(this.data)
.ruleChainId(this.ruleChainId)
.ruleNodeId(this.ruleNodeId)
.correlationId(this.correlationId)
.partition(this.partition)
.ctx(this.ctx)
.callback(this.callback);
return new TbMsgTransformer(this);
}
public TbMsgBuilder copy() {
return new TbMsgBuilder(this);
}
private static class TbMsgTransformer extends TbMsgBuilder {
TbMsgTransformer(TbMsg tbMsg) {
super(tbMsg);
}
/*
* metadata is only copied if specified explicitly during transform
* */
@Override
public TbMsgTransformer metaData(TbMsgMetaData metaData) {
super.metaData(metaData.copy());
this.metaData = metaData.copy();
return this;
}
/*
* setting ruleNodeId to null when updating ruleChainId
* */
@Override
public TbMsgTransformer ctx(TbMsgProcessingCtx ctx) {
super.ctx(ctx.copy());
public TbMsgBuilder ruleChainId(RuleChainId ruleChainId) {
this.ruleChainId = ruleChainId;
this.ruleNodeId = null;
return this;
}
@Override
public TbMsg build() {
/*
* always copying ctx when transforming
* */
if (ctx != null) {
ctx = ctx.copy();
}
return super.build();
}
}
private static class TbMsgBuilder {
private String queueName;
private UUID id;
private long ts;
private String type;
private TbMsgType internalType;
private EntityId originator;
private CustomerId customerId;
private TbMsgMetaData metaData;
private TbMsgDataType dataType;
private String data;
private RuleChainId ruleChainId;
private RuleNodeId ruleNodeId;
private UUID correlationId;
private Integer partition;
private TbMsgProcessingCtx ctx;
private TbMsgCallback callback;
public static class TbMsgBuilder {
protected String queueName;
protected UUID id;
protected long ts;
protected String type;
protected TbMsgType internalType;
protected EntityId originator;
protected CustomerId customerId;
protected TbMsgMetaData metaData;
protected TbMsgDataType dataType;
protected String data;
protected RuleChainId ruleChainId;
protected RuleNodeId ruleNodeId;
protected UUID correlationId;
protected Integer partition;
protected TbMsgProcessingCtx ctx;
protected TbMsgCallback callback;
TbMsgBuilder() {}
TbMsgBuilder(TbMsg tbMsg) {
this.queueName = tbMsg.queueName;
this.id = tbMsg.id;
this.ts = tbMsg.ts;
this.type = tbMsg.type;
this.internalType = tbMsg.internalType;
this.originator = tbMsg.originator;
this.customerId = tbMsg.customerId;
this.metaData = tbMsg.metaData;
this.dataType = tbMsg.dataType;
this.data = tbMsg.data;
this.ruleChainId = tbMsg.ruleChainId;
this.ruleNodeId = tbMsg.ruleNodeId;
this.correlationId = tbMsg.correlationId;
this.partition = tbMsg.partition;
this.ctx = tbMsg.ctx;
this.callback = tbMsg.callback;
}
public TbMsgBuilder queueName(String queueName) {
this.queueName = queueName;
return this;
@ -442,6 +412,7 @@ public final class TbMsg implements Serializable {
public TbMsgBuilder type(TbMsgType internalType) {
this.internalType = internalType;
this.type = internalType.name();
return this;
}

4
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbAbstractAlarmNode.java

@ -79,7 +79,9 @@ public abstract class TbAbstractAlarmNode<C extends TbAbstractAlarmNodeConfigura
if (previousDetails != null) {
TbMsgMetaData metaData = msg.getMetaData().copy();
metaData.putValue(PREV_ALARM_DETAILS, JacksonUtil.toString(previousDetails));
dummyMsg = TbMsg.transformMsgMetadata(msg, metaData);
dummyMsg = msg.transform()
.metaData(metaData)
.build();
}
return scriptEngine.executeJsonAsync(dummyMsg);
} catch (Exception e) {

9
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/lambda/TbAwsLambdaNode.java

@ -132,14 +132,19 @@ public class TbAwsLambdaNode extends TbAbstractExternalNode {
TbMsgMetaData metaData = originalMsg.getMetaData().copy();
metaData.putValue("requestId", invokeResult.getSdkResponseMetadata().getRequestId());
String data = getPayload(invokeResult);
return TbMsg.transformMsg(originalMsg, metaData, data);
return originalMsg.transform()
.metaData(metaData)
.data(data)
.build();
}
private TbMsg processException(TbMsg origMsg, InvokeResult invokeResult, Throwable t) {
TbMsgMetaData metaData = origMsg.getMetaData().copy();
metaData.putValue("error", t.getClass() + ": " + t.getMessage());
metaData.putValue("requestId", invokeResult.getSdkResponseMetadata().getRequestId());
return TbMsg.transformMsgMetadata(origMsg, metaData);
return origMsg.transform()
.metaData(metaData)
.build();
}
@Override

8
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/sns/TbSnsNode.java

@ -103,13 +103,17 @@ public class TbSnsNode extends TbAbstractExternalNode {
TbMsgMetaData metaData = origMsg.getMetaData().copy();
metaData.putValue(MESSAGE_ID, result.getMessageId());
metaData.putValue(REQUEST_ID, result.getSdkResponseMetadata().getRequestId());
return TbMsg.transformMsgMetadata(origMsg, metaData);
return origMsg.transform()
.metaData(metaData)
.build();
}
private TbMsg processException(TbMsg origMsg, Throwable t) {
TbMsgMetaData metaData = origMsg.getMetaData().copy();
metaData.putValue(ERROR, t.getClass() + ": " + t.getMessage());
return TbMsg.transformMsgMetadata(origMsg, metaData);
return origMsg.transform()
.metaData(metaData)
.build();
}
@Override

8
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/sqs/TbSqsNode.java

@ -134,13 +134,17 @@ public class TbSqsNode extends TbAbstractExternalNode {
if (!StringUtils.isEmpty(result.getSequenceNumber())) {
metaData.putValue(SEQUENCE_NUMBER, result.getSequenceNumber());
}
return TbMsg.transformMsgMetadata(origMsg, metaData);
return origMsg.transform()
.metaData(metaData)
.build();
}
private TbMsg processException(TbMsg origMsg, Throwable t) {
TbMsgMetaData metaData = origMsg.getMetaData().copy();
metaData.putValue(ERROR, t.getClass() + ": " + t.getMessage());
return TbMsg.transformMsgMetadata(origMsg, metaData);
return origMsg.transform()
.metaData(metaData)
.build();
}
@Override

8
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/gcp/pubsub/TbPubSubNode.java

@ -120,13 +120,17 @@ public class TbPubSubNode extends TbAbstractExternalNode {
private TbMsg processPublishResult(TbMsg origMsg, String messageId) {
TbMsgMetaData metaData = origMsg.getMetaData().copy();
metaData.putValue(MESSAGE_ID, messageId);
return TbMsg.transformMsgMetadata(origMsg, metaData);
return origMsg.transform()
.metaData(metaData)
.build();
}
private TbMsg processException(TbMsg origMsg, Throwable t) {
TbMsgMetaData metaData = origMsg.getMetaData().copy();
metaData.putValue(ERROR, t.getClass() + ": " + t.getMessage());
return TbMsg.transformMsgMetadata(origMsg, metaData);
return origMsg.transform()
.metaData(metaData)
.build();
}
Publisher initPubSubClient(TbContext ctx) throws IOException {

8
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/kafka/TbKafkaNode.java

@ -188,13 +188,17 @@ public class TbKafkaNode extends TbAbstractExternalNode {
metaData.putValue(OFFSET, String.valueOf(recordMetadata.offset()));
metaData.putValue(PARTITION, String.valueOf(recordMetadata.partition()));
metaData.putValue(TOPIC, recordMetadata.topic());
return TbMsg.transformMsgMetadata(origMsg, metaData);
return origMsg.transform()
.metaData(metaData)
.build();
}
private TbMsg processException(TbMsg origMsg, Exception e) {
TbMsgMetaData metaData = origMsg.getMetaData().copy();
metaData.putValue(ERROR, e.getClass() + ": " + e.getMessage());
return TbMsg.transformMsgMetadata(origMsg, metaData);
return origMsg.transform()
.metaData(metaData)
.build();
}
}

8
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/math/TbMathNode.java

@ -202,7 +202,9 @@ public class TbMathNode implements TbNode {
} else {
body.put(mathResultKey, toDoubleValue(mathResultDef, result));
}
return TbMsg.transformMsgData(msg, JacksonUtil.toString(body));
return msg.transform()
.data(JacksonUtil.toString(body))
.build();
}
private TbMsg addToMeta(TbMsg msg, TbMathResult mathResultDef, String mathResultKey, double result) {
@ -212,7 +214,9 @@ public class TbMathNode implements TbNode {
} else {
md.putValue(mathResultKey, Double.toString(toDoubleValue(mathResultDef, result)));
}
return TbMsg.transformMsgMetadata(msg, md);
return msg.transform()
.metaData(md)
.build();
}
private double calculateResult(List<TbMathArgumentValue> args) {

4
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/CalculateDeltaNode.java

@ -175,7 +175,9 @@ public class CalculateDeltaNode implements TbNode {
long period = previousData != null ? msg.getMetaDataTs() - previousData.ts : 0;
json.put(config.getPeriodValueKey(), period);
}
return TbMsg.transformMsgData(msg, JacksonUtil.toString(json));
return msg.transform()
.data(JacksonUtil.toString(json))
.build();
}, MoreExecutors.directExecutor());
}

8
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractNodeWithFetchTo.java

@ -82,9 +82,13 @@ public abstract class TbAbstractNodeWithFetchTo<C extends TbAbstractFetchToNodeC
protected TbMsg transformMessage(TbMsg msg, ObjectNode msgDataNode, TbMsgMetaData msgMetaData) {
switch (fetchTo) {
case DATA:
return TbMsg.transformMsgData(msg, JacksonUtil.toString(msgDataNode));
return msg.transform()
.data(JacksonUtil.toString(msgDataNode))
.build();
case METADATA:
return TbMsg.transformMsgMetadata(msg, msgMetaData);
return msg.transform()
.metaData(msgMetaData)
.build();
default:
log.debug("Unexpected FetchTo value: {}. Allowed values: {}", fetchTo, TbMsgSource.values());
return msg;

4
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNode.java

@ -115,7 +115,9 @@ public class TbGetTelemetryNode implements TbNode {
ListenableFuture<List<TsKvEntry>> list = ctx.getTimeseriesService().findAll(ctx.getTenantId(), msg.getOriginator(), buildQueries(interval, keys));
DonAsynchron.withCallback(list, data -> {
var metaData = updateMetadata(data, msg, keys);
ctx.tellSuccess(TbMsg.transformMsgMetadata(msg, metaData));
ctx.tellSuccess(msg.transform()
.metaData(metaData)
.build());
}, error -> ctx.tellFailure(msg, error), ctx.getDbCallbackExecutor());
}

4
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mqtt/TbMqttNode.java

@ -103,7 +103,9 @@ public class TbMqttNode extends TbAbstractExternalNode {
private TbMsg processException(TbMsg origMsg, Throwable e) {
TbMsgMetaData metaData = origMsg.getMetaData().copy();
metaData.putValue(ERROR, e.getClass() + ": " + e.getMessage());
return TbMsg.transformMsgMetadata(origMsg, metaData);
return origMsg.transform()
.metaData(metaData)
.build();
}
@Override

4
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbNotificationNode.java

@ -82,7 +82,9 @@ public class TbNotificationNode extends TbAbstractExternalNode {
public void onSuccess(NotificationRequestStats stats) {
TbMsgMetaData metaData = tbMsg.getMetaData().copy();
metaData.putValue("notificationRequestResult", JacksonUtil.toString(stats));
tellSuccess(ctx, TbMsg.transformMsgMetadata(tbMsg, metaData));
tellSuccess(ctx, tbMsg.transform()
.metaData(metaData)
.build());
}
@Override

4
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rabbitmq/TbRabbitMqNode.java

@ -127,7 +127,9 @@ public class TbRabbitMqNode extends TbAbstractExternalNode {
private TbMsg processException(TbMsg origMsg, Throwable t) {
TbMsgMetaData metaData = origMsg.getMetaData().copy();
metaData.putValue(ERROR, t.getClass() + ": " + t.getMessage());
return TbMsg.transformMsgMetadata(origMsg, metaData);
return origMsg.transform()
.metaData(metaData)
.build();
}
@Override

8
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbHttpClient.java

@ -323,7 +323,9 @@ public class TbHttpClient {
metaData.putValue(STATUS_REASON, httpStatus.getReasonPhrase());
metaData.putValue(ERROR_BODY, response.getBody());
headersToMetaData(response.getHeaders(), metaData::putValue);
return TbMsg.transformMsgMetadata(origMsg, metaData);
return origMsg.transform()
.metaData(metaData)
.build();
}
private TbMsg processException(TbMsg origMsg, Throwable e) {
@ -334,7 +336,9 @@ public class TbHttpClient {
metaData.putValue(STATUS_CODE, restClientResponseException.getStatusCode().value() + "");
metaData.putValue(ERROR_BODY, restClientResponseException.getResponseBodyAsString());
}
return TbMsg.transformMsgMetadata(origMsg, metaData);
return origMsg.transform()
.metaData(metaData)
.build();
}
private void prepareHeaders(HttpHeaders headers, TbMsg msg) {

5
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbCopyKeysNode.java

@ -105,7 +105,10 @@ public class TbCopyKeysNode extends TbAbstractTransformNodeWithTbMsgSource {
log.debug("Unexpected CopyFrom value: {}. Allowed values: {}", copyFrom, TbMsgSource.values());
}
}
ctx.tellSuccess(msgChanged ? TbMsg.transformMsg(msg, metaDataCopy, msgData) : msg);
ctx.tellSuccess(msgChanged ? msg.transform()
.metaData(metaDataCopy)
.data(msgData)
.build() : msg);
}
@Override

5
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbDeleteKeysNode.java

@ -100,7 +100,10 @@ public class TbDeleteKeysNode extends TbAbstractTransformNodeWithTbMsgSource {
default:
log.debug("Unexpected DeleteFrom value: {}. Allowed values: {}", deleteFrom, TbMsgSource.values());
}
ctx.tellSuccess(hasNoChanges ? msg : TbMsg.transformMsg(msg, metaDataCopy, msgDataStr));
ctx.tellSuccess(hasNoChanges ? msg : msg.transform()
.metaData(metaDataCopy)
.data(msgDataStr)
.build());
}
@Override

4
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbJsonPathNode.java

@ -68,7 +68,9 @@ public class TbJsonPathNode implements TbNode {
if (!TbJsonPathNodeConfiguration.DEFAULT_JSON_PATH.equals(this.jsonPathValue)) {
try {
Object jsonPathData = jsonPath.read(msg.getData(), this.configurationJsonPath);
ctx.tellSuccess(TbMsg.transformMsgData(msg, JacksonUtil.toString(jsonPathData)));
ctx.tellSuccess(msg.transform()
.data(JacksonUtil.toString(jsonPathData))
.build());
} catch (PathNotFoundException e) {
ctx.tellFailure(msg, e);
}

5
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbRenameKeysNode.java

@ -106,7 +106,10 @@ public class TbRenameKeysNode extends TbAbstractTransformNodeWithTbMsgSource {
default:
log.debug("Unexpected RenameIn value: {}. Allowed values: {}", renameIn, TbMsgSource.values());
}
ctx.tellSuccess(msgChanged ? TbMsg.transformMsg(msg, metaDataCopy, data) : msg);
ctx.tellSuccess(msgChanged ? msg.transform()
.metaData(metaDataCopy)
.data(data)
.build() : msg);
}
@Override

8
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbSplitArrayMsgNode.java

@ -64,7 +64,9 @@ public class TbSplitArrayMsgNode implements TbNode {
if (data.isEmpty()) {
ctx.ack(msg);
} else if (data.size() == 1) {
ctx.tellSuccess(TbMsg.transformMsgData(msg, JacksonUtil.toString(data.get(0))));
ctx.tellSuccess(msg.transform()
.data(JacksonUtil.toString(data.get(0)))
.build());
} else {
TbMsgCallbackWrapper wrapper = new MultipleTbMsgsCallbackWrapper(data.size(), new TbMsgCallback() {
@Override
@ -78,7 +80,9 @@ public class TbSplitArrayMsgNode implements TbNode {
}
});
data.forEach(msgNode -> {
TbMsg outMsg = TbMsg.transformMsgData(msg, JacksonUtil.toString(msgNode));
TbMsg outMsg = msg.transform()
.data(JacksonUtil.toString(msgNode))
.build();
ctx.enqueueForTellNext(outMsg, TbNodeConnectionType.SUCCESS, wrapper::onSuccess, wrapper::onFailure);
});
}

72
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/action/TbCreateAlarmNodeTest.java

@ -229,12 +229,12 @@ class TbCreateAlarmNodeTest {
given(alarmServiceMock.createAlarm(expectedCreateAlarmRequest)).willReturn(apiCallResult);
given(ctxMock.alarmActionMsg(expectedCreatedAlarmInfo, ruleNodeSelfId, TbMsgType.ENTITY_CREATED)).willReturn(alarmActionMsgMock);
given(ctxMock.transformMsg(any(TbMsg.class), any(TbMsgType.class), any(EntityId.class), any(TbMsgMetaData.class), anyString()))
.willAnswer(answer -> TbMsg.transformMsg(
answer.getArgument(0, TbMsg.class),
answer.getArgument(1, TbMsgType.class),
answer.getArgument(2, EntityId.class),
answer.getArgument(3, TbMsgMetaData.class),
answer.getArgument(4, String.class))
.willAnswer(answer -> answer.getArgument(0, TbMsg.class).transform()
.type(answer.getArgument(1, TbMsgType.class))
.originator(answer.getArgument(2, EntityId.class))
.metaData(answer.getArgument(3, TbMsgMetaData.class))
.data(answer.getArgument(4, String.class))
.build()
);
given(ctxMock.createScriptEngine(ScriptLanguage.TBEL, TbAbstractAlarmNodeConfiguration.ALARM_DETAILS_BUILD_TBEL_TEMPLATE)).willReturn(alarmDetailsScriptMock);
@ -403,12 +403,12 @@ class TbCreateAlarmNodeTest {
given(alarmServiceMock.createAlarm(expectedCreateAlarmRequest)).willReturn(apiCallResult);
given(ctxMock.alarmActionMsg(expectedCreatedAlarmInfo, ruleNodeSelfId, TbMsgType.ENTITY_CREATED)).willReturn(alarmActionMsgMock);
given(ctxMock.transformMsg(any(TbMsg.class), any(TbMsgType.class), any(EntityId.class), any(TbMsgMetaData.class), anyString()))
.willAnswer(answer -> TbMsg.transformMsg(
answer.getArgument(0, TbMsg.class),
answer.getArgument(1, TbMsgType.class),
answer.getArgument(2, EntityId.class),
answer.getArgument(3, TbMsgMetaData.class),
answer.getArgument(4, String.class))
.willAnswer(answer -> answer.getArgument(0, TbMsg.class).transform()
.type(answer.getArgument(1, TbMsgType.class))
.originator(answer.getArgument(2, EntityId.class))
.metaData(answer.getArgument(3, TbMsgMetaData.class))
.data(answer.getArgument(4, String.class))
.build()
);
given(ctxMock.createScriptEngine(ScriptLanguage.JS, config.getAlarmDetailsBuildJs())).willReturn(alarmDetailsScriptMock);
@ -598,12 +598,12 @@ class TbCreateAlarmNodeTest {
given(alarmServiceMock.updateAlarm(expectedUpdateAlarmRequest)).willReturn(apiCallResult);
given(ctxMock.alarmActionMsg(expectedUpdatedAlarmInfo, ruleNodeSelfId, TbMsgType.ENTITY_UPDATED)).willReturn(alarmActionMsgMock);
given(ctxMock.transformMsg(any(TbMsg.class), any(TbMsgType.class), any(EntityId.class), any(TbMsgMetaData.class), anyString()))
.willAnswer(answer -> TbMsg.transformMsg(
answer.getArgument(0, TbMsg.class),
answer.getArgument(1, TbMsgType.class),
answer.getArgument(2, EntityId.class),
answer.getArgument(3, TbMsgMetaData.class),
answer.getArgument(4, String.class))
.willAnswer(answer -> answer.getArgument(0, TbMsg.class).transform()
.type(answer.getArgument(1, TbMsgType.class))
.originator(answer.getArgument(2, EntityId.class))
.metaData(answer.getArgument(3, TbMsgMetaData.class))
.data(answer.getArgument(4, String.class))
.build()
);
given(ctxMock.createScriptEngine(ScriptLanguage.TBEL, config.getAlarmDetailsBuildTbel())).willReturn(alarmDetailsScriptMock);
@ -775,12 +775,12 @@ class TbCreateAlarmNodeTest {
given(alarmServiceMock.createAlarm(expectedCreateAlarmRequest)).willReturn(apiCallResult);
given(ctxMock.alarmActionMsg(expectedCreatedAlarmInfo, ruleNodeSelfId, TbMsgType.ENTITY_CREATED)).willReturn(alarmActionMsgMock);
given(ctxMock.transformMsg(any(TbMsg.class), any(TbMsgType.class), any(EntityId.class), any(TbMsgMetaData.class), anyString()))
.willAnswer(answer -> TbMsg.transformMsg(
answer.getArgument(0, TbMsg.class),
answer.getArgument(1, TbMsgType.class),
answer.getArgument(2, EntityId.class),
answer.getArgument(3, TbMsgMetaData.class),
answer.getArgument(4, String.class))
.willAnswer(answer -> answer.getArgument(0, TbMsg.class).transform()
.type(answer.getArgument(1, TbMsgType.class))
.originator(answer.getArgument(2, EntityId.class))
.metaData(answer.getArgument(3, TbMsgMetaData.class))
.data(answer.getArgument(4, String.class))
.build()
);
given(ctxMock.createScriptEngine(ScriptLanguage.TBEL, config.getAlarmDetailsBuildTbel())).willReturn(alarmDetailsScriptMock);
@ -967,12 +967,12 @@ class TbCreateAlarmNodeTest {
given(alarmServiceMock.updateAlarm(expectedUpdateAlarmRequest)).willReturn(apiCallResult);
given(ctxMock.alarmActionMsg(expectedUpdatedAlarmInfo, ruleNodeSelfId, TbMsgType.ENTITY_UPDATED)).willReturn(alarmActionMsgMock);
given(ctxMock.transformMsg(any(TbMsg.class), any(TbMsgType.class), any(EntityId.class), any(TbMsgMetaData.class), anyString()))
.willAnswer(answer -> TbMsg.transformMsg(
answer.getArgument(0, TbMsg.class),
answer.getArgument(1, TbMsgType.class),
answer.getArgument(2, EntityId.class),
answer.getArgument(3, TbMsgMetaData.class),
answer.getArgument(4, String.class))
.willAnswer(answer -> answer.getArgument(0, TbMsg.class).transform()
.type(answer.getArgument(1, TbMsgType.class))
.originator(answer.getArgument(2, EntityId.class))
.metaData(answer.getArgument(3, TbMsgMetaData.class))
.data(answer.getArgument(4, String.class))
.build()
);
given(ctxMock.createScriptEngine(ScriptLanguage.TBEL, config.getAlarmDetailsBuildTbel())).willReturn(alarmDetailsScriptMock);
@ -1153,12 +1153,12 @@ class TbCreateAlarmNodeTest {
given(alarmServiceMock.updateAlarm(expectedUpdateAlarmRequest)).willReturn(apiCallResult);
given(ctxMock.alarmActionMsg(expectedUpdatedAlarmInfo, ruleNodeSelfId, TbMsgType.ENTITY_UPDATED)).willReturn(alarmActionMsgMock);
given(ctxMock.transformMsg(any(TbMsg.class), any(TbMsgType.class), any(EntityId.class), any(TbMsgMetaData.class), anyString()))
.willAnswer(answer -> TbMsg.transformMsg(
answer.getArgument(0, TbMsg.class),
answer.getArgument(1, TbMsgType.class),
answer.getArgument(2, EntityId.class),
answer.getArgument(3, TbMsgMetaData.class),
answer.getArgument(4, String.class))
.willAnswer(answer -> answer.getArgument(0, TbMsg.class).transform()
.type(answer.getArgument(1, TbMsgType.class))
.originator(answer.getArgument(2, EntityId.class))
.metaData(answer.getArgument(3, TbMsgMetaData.class))
.data(answer.getArgument(4, String.class))
.build()
);
given(ctxMock.createScriptEngine(ScriptLanguage.TBEL, config.getAlarmDetailsBuildTbel())).willReturn(alarmDetailsScriptMock);

4
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/action/TbCreateRelationNodeTest.java

@ -439,7 +439,9 @@ public class TbCreateRelationNodeTest extends AbstractRuleNodeUpgradeTest {
var md = getMetadataWithNameTemplate();
var msg = getTbMsg(originatorId, md);
var msgAfterOriginatorChanged = TbMsg.transformMsgOriginator(msg, originatorId);
var msgAfterOriginatorChanged = msg.transform()
.originator(originatorId)
.build();
when(ctxMock.transformMsgOriginator(any(), any())).thenReturn(msgAfterOriginatorChanged);
// WHEN

18
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/aws/lambda/TbAwsLambdaNodeTest.java

@ -176,7 +176,10 @@ public class TbAwsLambdaNodeTest {
assertThat(invokeRequestCaptor.getValue().getQualifier()).isEqualTo(expectedQualifier);
TbMsgMetaData resultMsgMetadata = metadata.copy();
resultMsgMetadata.putValue("requestId", requestIdStr);
TbMsg resultedMsg = TbMsg.transformMsg(msg, resultMsgMetadata, funcResponsePayload);
TbMsg resultedMsg = msg.transform()
.metaData(resultMsgMetadata)
.data(funcResponsePayload)
.build();
assertThat(msgCaptor.getValue()).usingRecursiveComparison()
.ignoringFields("ctx")
.isEqualTo(resultedMsg);
@ -231,7 +234,9 @@ public class TbAwsLambdaNodeTest {
verify(ctx).tellFailure(msgCaptor.capture(), throwableCaptor.capture());
var metadata = Map.of("error", RuntimeException.class + ": " + errorMsg, "requestId", requestIdStr);
TbMsg resultedMsg = TbMsg.transformMsgMetadata(msg, new TbMsgMetaData(metadata));
TbMsg resultedMsg = msg.transform()
.metaData(new TbMsgMetaData(metadata))
.build();
assertThat(msgCaptor.getValue()).usingRecursiveComparison()
.ignoringFields("ctx")
@ -270,7 +275,10 @@ public class TbAwsLambdaNodeTest {
verify(ctx).tellSuccess(msgCaptor.capture());
Map<String, String> metadata = Map.of("requestId", requestIdStr);
TbMsg resultedMsg = TbMsg.transformMsg(msg, new TbMsgMetaData(metadata), payload);
TbMsg resultedMsg = msg.transform()
.metaData(new TbMsgMetaData(metadata))
.data(payload)
.build();
assertThat(msgCaptor.getValue()).usingRecursiveComparison()
.ignoringFields("ctx")
@ -309,7 +317,9 @@ public class TbAwsLambdaNodeTest {
verify(ctx).tellFailure(msgCaptor.capture(), throwableCaptor.capture());
var metadata = Map.of("error", RuntimeException.class + ": " + errorMsg, "requestId", requestIdStr);
TbMsg resultedMsg = TbMsg.transformMsgMetadata(msg, new TbMsgMetaData(metadata));
TbMsg resultedMsg = msg.transform()
.metaData(new TbMsgMetaData(metadata))
.build();
assertThat(msgCaptor.getValue()).usingRecursiveComparison()
.ignoringFields("ctx")

16
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/gcp/pubsub/TbPubSubNodeTest.java

@ -138,7 +138,9 @@ class TbPubSubNodeTest {
ArgumentCaptor<TbMsg> actualMsg = ArgumentCaptor.forClass(TbMsg.class);
then(ctxMock).should().enqueueForTellNext(actualMsg.capture(), eq(TbNodeConnectionType.SUCCESS));
metaData.putValue("messageId", messageId);
TbMsg expectedMsg = TbMsg.transformMsgMetadata(msg, metaData);
TbMsg expectedMsg = msg.transform()
.metaData(metaData)
.build();
assertThat(actualMsg.getValue())
.usingRecursiveComparison()
.ignoringFields("ctx")
@ -184,7 +186,9 @@ class TbPubSubNodeTest {
ArgumentCaptor<TbMsg> actualMsg = ArgumentCaptor.forClass(TbMsg.class);
then(ctxMock).should().tellSuccess(actualMsg.capture());
metadata.putValue("messageId", messageId);
TbMsg expectedMsg = TbMsg.transformMsgMetadata(msg, metadata);
TbMsg expectedMsg = msg.transform()
.metaData(metadata)
.build();
assertThat(actualMsg.getValue())
.usingRecursiveComparison()
.ignoringFields("ctx")
@ -216,7 +220,9 @@ class TbPubSubNodeTest {
ArgumentCaptor<Throwable> actualError = ArgumentCaptor.forClass(Throwable.class);
then(ctxMock).should().tellFailure(actualMsg.capture(), actualError.capture());
metaData.putValue("error", RuntimeException.class + ": " + errorMsg);
TbMsg expectedMsg = TbMsg.transformMsgMetadata(msg, metaData);
TbMsg expectedMsg = msg.transform()
.metaData(metaData)
.build();
assertThat(actualMsg.getValue())
.usingRecursiveComparison()
.ignoringFields("ctx")
@ -249,7 +255,9 @@ class TbPubSubNodeTest {
ArgumentCaptor<Throwable> actualError = ArgumentCaptor.forClass(Throwable.class);
then(ctxMock).should().enqueueForTellFailure(actualMsg.capture(), actualError.capture());
metaData.putValue("error", RuntimeException.class + ": " + errorMsg);
TbMsg expectedMsg = TbMsg.transformMsgMetadata(msg, metaData);
TbMsg expectedMsg = msg.transform()
.metaData(metaData)
.build();
assertThat(actualMsg.getValue())
.usingRecursiveComparison()
.ignoringFields("ctx")

8
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/kafka/TbKafkaNodeTest.java

@ -436,7 +436,9 @@ public class TbKafkaNodeTest {
metaData.putValue("offset", String.valueOf(OFFSET));
metaData.putValue("partition", String.valueOf(PARTITION));
metaData.putValue("topic", expectedTopic);
TbMsg expectedMsg = TbMsg.transformMsgMetadata(originalMsg, metaData);
TbMsg expectedMsg = originalMsg.transform()
.metaData(metaData)
.build();
assertThat(actualMsg)
.usingRecursiveComparison()
.ignoringFields("ctx")
@ -446,7 +448,9 @@ public class TbKafkaNodeTest {
private void verifyOutgoingFailureMsg(String errorMsg, TbMsg actualMsg, TbMsg originalMsg) {
TbMsgMetaData metaData = originalMsg.getMetaData();
metaData.putValue("error", RuntimeException.class + ": " + errorMsg);
TbMsg expectedMsg = TbMsg.transformMsgMetadata(originalMsg, metaData);
TbMsg expectedMsg = originalMsg.transform()
.metaData(metaData)
.build();
assertThat(actualMsg).usingRecursiveComparison().ignoringFields("ctx").isEqualTo(expectedMsg);
}

8
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNodeTest.java

@ -465,7 +465,9 @@ public class TbGetTelemetryNodeTest extends AbstractRuleNodeUpgradeTest {
TbMsgMetaData metaData = new TbMsgMetaData();
metaData.putValue("temperature", "[{\"ts\":" + (ts - 5) + ",\"value\":23.1},{\"ts\":" + (ts - 4) + ",\"value\":22.4}]");
metaData.putValue("humidity", "[{\"ts\":" + (ts - 4) + ",\"value\":55.5}]");
TbMsg expectedMsg = TbMsg.transformMsgMetadata(msg, metaData);
TbMsg expectedMsg = msg.transform()
.metaData(metaData)
.build();
assertThat(actualMsg.getValue()).usingRecursiveComparison().ignoringFields("ctx").isEqualTo(expectedMsg);
}
@ -501,7 +503,9 @@ public class TbGetTelemetryNodeTest extends AbstractRuleNodeUpgradeTest {
TbMsgMetaData metaData = new TbMsgMetaData();
metaData.putValue("temperature", "\"22.4\"");
metaData.putValue("humidity", "\"55.5\"");
TbMsg expectedMsg = TbMsg.transformMsgMetadata(msg, metaData);
TbMsg expectedMsg = msg.transform()
.metaData(metaData)
.build();
assertThat(actualMsg.getValue()).usingRecursiveComparison().ignoringFields("ctx").isEqualTo(expectedMsg);
}

4
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/mqtt/TbMqttNodeTest.java

@ -340,7 +340,9 @@ public class TbMqttNodeTest extends AbstractRuleNodeUpgradeTest {
then(mqttClientMock).should().publish(mqttNodeConfig.getTopicPattern(), Unpooled.wrappedBuffer(expectedData.getBytes(StandardCharsets.UTF_8)), MqttQoS.AT_LEAST_ONCE, false);
TbMsgMetaData metaData = new TbMsgMetaData();
metaData.putValue("error", RuntimeException.class + ": " + errorMsg);
TbMsg expectedMsg = TbMsg.transformMsgMetadata(msg, metaData);
TbMsg expectedMsg = msg.transform()
.metaData(metaData)
.build();
ArgumentCaptor<TbMsg> actualMsgCaptor = ArgumentCaptor.forClass(TbMsg.class);
then(ctxMock).should().tellFailure(actualMsgCaptor.capture(), eq(exception));
TbMsg actualMsg = actualMsgCaptor.getValue();

4
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/rabbitmq/TbRabbitMqNodeTest.java

@ -227,7 +227,9 @@ public class TbRabbitMqNodeTest {
() -> then(ctxMock).should().tellFailure(actualMsg.capture(), throwable.capture());
verifyTellFailure.run();
metaData.putValue("error", RuntimeException.class + ": " + errorMsg);
TbMsg expectedMsg = TbMsg.transformMsgMetadata(msg, metaData);
TbMsg expectedMsg = msg.transform()
.metaData(metaData)
.build();
assertThat(actualMsg.getValue()).usingRecursiveComparison().ignoringFields("ctx").isEqualTo(expectedMsg);
assertThat(throwable.getValue()).isInstanceOf(RuntimeException.class).hasMessage(errorMsg);
}

16
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNodeTest.java

@ -182,7 +182,9 @@ public class TbChangeOriginatorNodeTest {
.metaData(TbMsgMetaData.EMPTY.copy())
.data(TbMsg.EMPTY_JSON_OBJECT)
.build();
TbMsg expectedMsg = TbMsg.transformMsgOriginator(msg, CUSTOMER_ID);
TbMsg expectedMsg = msg.transform()
.originator(CUSTOMER_ID)
.build();
given(ctxMock.getDbCallbackExecutor()).willReturn(dbExecutor);
given(ctxMock.getDeviceService()).willReturn(deviceServiceMock);
@ -210,7 +212,9 @@ public class TbChangeOriginatorNodeTest {
.metaData(TbMsgMetaData.EMPTY.copy())
.data(TbMsg.EMPTY_JSON_OBJECT)
.build();
TbMsg expectedMsg = TbMsg.transformMsgOriginator(msg, TENANT_ID);
TbMsg expectedMsg = msg.transform()
.originator(TENANT_ID)
.build();
given(ctxMock.getDbCallbackExecutor()).willReturn(dbExecutor);
given(ctxMock.getTenantId()).willReturn(TENANT_ID);
@ -274,7 +278,9 @@ public class TbChangeOriginatorNodeTest {
.metaData(TbMsgMetaData.EMPTY.copy())
.data(TbMsg.EMPTY_JSON_OBJECT)
.build();
TbMsg expectedMsg = TbMsg.transformMsgOriginator(msg, DEVICE_ID);
TbMsg expectedMsg = msg.transform()
.originator(DEVICE_ID)
.build();
given(ctxMock.getDbCallbackExecutor()).willReturn(dbExecutor);
given(ctxMock.getAlarmService()).willReturn(alarmServiceMock);
@ -305,7 +311,9 @@ public class TbChangeOriginatorNodeTest {
.metaData(metaData.copy())
.data(data)
.build();
TbMsg expectedMsg = TbMsg.transformMsgOriginator(msg, ASSET_ID);
TbMsg expectedMsg = msg.transform()
.originator(ASSET_ID)
.build();
given(ctxMock.getDbCallbackExecutor()).willReturn(dbExecutor);
given(ctxMock.getAssetService()).willReturn(assetServiceMock);

Loading…
Cancel
Save