|
|
|
@ -52,6 +52,7 @@ public final class TbMsg implements Serializable { |
|
|
|
private final UUID id; |
|
|
|
private final long ts; |
|
|
|
private final String type; |
|
|
|
private final TbMsgType internalType; |
|
|
|
private final EntityId originator; |
|
|
|
private final CustomerId customerId; |
|
|
|
private final TbMsgMetaData metaData; |
|
|
|
@ -97,7 +98,7 @@ public final class TbMsg implements Serializable { |
|
|
|
*/ |
|
|
|
@Deprecated(since = "3.5.2") |
|
|
|
public static TbMsg newMsg(String queueName, String type, EntityId originator, CustomerId customerId, TbMsgMetaData metaData, String data, RuleChainId ruleChainId, RuleNodeId ruleNodeId) { |
|
|
|
return new TbMsg(queueName, UUID.randomUUID(), System.currentTimeMillis(), type, originator, customerId, |
|
|
|
return new TbMsg(queueName, UUID.randomUUID(), System.currentTimeMillis(), null, type, originator, customerId, |
|
|
|
metaData.copy(), TbMsgDataType.JSON, data, ruleChainId, ruleNodeId, null, TbMsgCallback.EMPTY); |
|
|
|
} |
|
|
|
|
|
|
|
@ -108,7 +109,7 @@ public final class TbMsg implements Serializable { |
|
|
|
|
|
|
|
@Deprecated(since = "3.5.2", forRemoval = true) |
|
|
|
public static TbMsg newMsg(String type, EntityId originator, CustomerId customerId, TbMsgMetaData metaData, String data) { |
|
|
|
return new TbMsg(null, UUID.randomUUID(), System.currentTimeMillis(), type, originator, customerId, |
|
|
|
return new TbMsg(null, UUID.randomUUID(), System.currentTimeMillis(), null, type, originator, customerId, |
|
|
|
metaData.copy(), TbMsgDataType.JSON, data, null, null, null, TbMsgCallback.EMPTY); |
|
|
|
} |
|
|
|
|
|
|
|
@ -117,7 +118,7 @@ public final class TbMsg implements Serializable { |
|
|
|
} |
|
|
|
|
|
|
|
public static TbMsg newMsg(String queueName, TbMsgType type, EntityId originator, CustomerId customerId, TbMsgMetaData metaData, String data, RuleChainId ruleChainId, RuleNodeId ruleNodeId) { |
|
|
|
return new TbMsg(queueName, UUID.randomUUID(), System.currentTimeMillis(), type.name(), originator, customerId, |
|
|
|
return new TbMsg(queueName, UUID.randomUUID(), System.currentTimeMillis(), type, originator, customerId, |
|
|
|
metaData.copy(), TbMsgDataType.JSON, data, ruleChainId, ruleNodeId, null, TbMsgCallback.EMPTY); |
|
|
|
} |
|
|
|
|
|
|
|
@ -126,7 +127,7 @@ public final class TbMsg implements Serializable { |
|
|
|
} |
|
|
|
|
|
|
|
public static TbMsg newMsg(TbMsgType type, EntityId originator, CustomerId customerId, TbMsgMetaData metaData, String data) { |
|
|
|
return new TbMsg(null, UUID.randomUUID(), System.currentTimeMillis(), type.name(), originator, customerId, |
|
|
|
return new TbMsg(null, UUID.randomUUID(), System.currentTimeMillis(), type, originator, customerId, |
|
|
|
metaData.copy(), TbMsgDataType.JSON, data, null, null, null, TbMsgCallback.EMPTY); |
|
|
|
} |
|
|
|
|
|
|
|
@ -170,13 +171,13 @@ public final class TbMsg implements Serializable { |
|
|
|
*/ |
|
|
|
@Deprecated(since = "3.5.2") |
|
|
|
public static TbMsg newMsg(String queueName, String type, EntityId originator, CustomerId customerId, TbMsgMetaData metaData, String data) { |
|
|
|
return new TbMsg(queueName, UUID.randomUUID(), System.currentTimeMillis(), type, originator, customerId, |
|
|
|
return new TbMsg(queueName, UUID.randomUUID(), System.currentTimeMillis(), null, type, originator, customerId, |
|
|
|
metaData.copy(), TbMsgDataType.JSON, data, null, null, null, TbMsgCallback.EMPTY); |
|
|
|
} |
|
|
|
|
|
|
|
@Deprecated(since = "3.5.2", forRemoval = true) |
|
|
|
public static TbMsg newMsg(String type, EntityId originator, CustomerId customerId, TbMsgMetaData metaData, TbMsgDataType dataType, String data) { |
|
|
|
return new TbMsg(null, UUID.randomUUID(), System.currentTimeMillis(), type, originator, customerId, |
|
|
|
return new TbMsg(null, UUID.randomUUID(), System.currentTimeMillis(), null, type, originator, customerId, |
|
|
|
metaData.copy(), dataType, data, null, null, null, TbMsgCallback.EMPTY); |
|
|
|
} |
|
|
|
|
|
|
|
@ -205,12 +206,12 @@ public final class TbMsg implements Serializable { |
|
|
|
} |
|
|
|
|
|
|
|
public static TbMsg newMsg(String queueName, TbMsgType type, EntityId originator, CustomerId customerId, TbMsgMetaData metaData, String data) { |
|
|
|
return new TbMsg(queueName, UUID.randomUUID(), System.currentTimeMillis(), type.name(), originator, customerId, |
|
|
|
return new TbMsg(queueName, UUID.randomUUID(), System.currentTimeMillis(), type, originator, customerId, |
|
|
|
metaData.copy(), TbMsgDataType.JSON, data, null, null, null, TbMsgCallback.EMPTY); |
|
|
|
} |
|
|
|
|
|
|
|
public static TbMsg newMsg(TbMsgType type, EntityId originator, CustomerId customerId, TbMsgMetaData metaData, TbMsgDataType dataType, String data) { |
|
|
|
return new TbMsg(null, UUID.randomUUID(), System.currentTimeMillis(), type.name(), originator, customerId, |
|
|
|
return new TbMsg(null, UUID.randomUUID(), System.currentTimeMillis(), type, originator, customerId, |
|
|
|
metaData.copy(), dataType, data, null, null, null, TbMsgCallback.EMPTY); |
|
|
|
} |
|
|
|
|
|
|
|
@ -222,13 +223,13 @@ public final class TbMsg implements Serializable { |
|
|
|
|
|
|
|
@Deprecated(since = "3.5.2", forRemoval = true) |
|
|
|
public static TbMsg newMsg(String type, EntityId originator, TbMsgMetaData metaData, TbMsgDataType dataType, String data, RuleChainId ruleChainId, RuleNodeId ruleNodeId) { |
|
|
|
return new TbMsg(null, UUID.randomUUID(), System.currentTimeMillis(), type, originator, null, |
|
|
|
return new TbMsg(null, UUID.randomUUID(), System.currentTimeMillis(), null, type, originator, null, |
|
|
|
metaData.copy(), dataType, data, ruleChainId, ruleNodeId, null, TbMsgCallback.EMPTY); |
|
|
|
} |
|
|
|
|
|
|
|
@Deprecated(since = "3.5.2", forRemoval = true) |
|
|
|
public static TbMsg newMsg(String type, EntityId originator, TbMsgMetaData metaData, String data, TbMsgCallback callback) { |
|
|
|
return new TbMsg(null, UUID.randomUUID(), System.currentTimeMillis(), type, originator, null, |
|
|
|
return new TbMsg(null, UUID.randomUUID(), System.currentTimeMillis(), null, type, originator, null, |
|
|
|
metaData.copy(), TbMsgDataType.JSON, data, null, null, null, callback); |
|
|
|
} |
|
|
|
|
|
|
|
@ -250,72 +251,77 @@ public final class TbMsg implements Serializable { |
|
|
|
*/ |
|
|
|
@Deprecated(since = "3.5.2") |
|
|
|
public static TbMsg transformMsg(TbMsg tbMsg, String type, EntityId originator, TbMsgMetaData metaData, String data) { |
|
|
|
return new TbMsg(tbMsg.queueName, tbMsg.id, tbMsg.ts, type, originator, tbMsg.customerId, metaData.copy(), tbMsg.dataType, |
|
|
|
return new TbMsg(tbMsg.queueName, tbMsg.id, tbMsg.ts, null, type, originator, tbMsg.customerId, metaData.copy(), tbMsg.dataType, |
|
|
|
data, tbMsg.ruleChainId, tbMsg.ruleNodeId, tbMsg.ctx.copy(), tbMsg.callback); |
|
|
|
} |
|
|
|
|
|
|
|
public static TbMsg newMsg(TbMsgType type, EntityId originator, TbMsgMetaData metaData, TbMsgDataType dataType, String data, RuleChainId ruleChainId, RuleNodeId ruleNodeId) { |
|
|
|
return new TbMsg(null, UUID.randomUUID(), System.currentTimeMillis(), type.name(), originator, null, |
|
|
|
return new TbMsg(null, UUID.randomUUID(), System.currentTimeMillis(), type, originator, null, |
|
|
|
metaData.copy(), dataType, data, ruleChainId, ruleNodeId, null, TbMsgCallback.EMPTY); |
|
|
|
} |
|
|
|
|
|
|
|
public static TbMsg newMsg(TbMsgType type, EntityId originator, TbMsgMetaData metaData, String data, TbMsgCallback callback) { |
|
|
|
return new TbMsg(null, UUID.randomUUID(), System.currentTimeMillis(), type.name(), originator, null, |
|
|
|
return new TbMsg(null, UUID.randomUUID(), System.currentTimeMillis(), type, originator, null, |
|
|
|
metaData.copy(), TbMsgDataType.JSON, data, null, null, null, callback); |
|
|
|
} |
|
|
|
|
|
|
|
public static TbMsg transformMsg(TbMsg tbMsg, TbMsgType type, EntityId originator, TbMsgMetaData metaData, String data) { |
|
|
|
return new TbMsg(tbMsg.queueName, tbMsg.id, tbMsg.ts, type.name(), originator, tbMsg.customerId, metaData.copy(), tbMsg.dataType, |
|
|
|
return new TbMsg(tbMsg.queueName, tbMsg.id, tbMsg.ts, type, originator, tbMsg.customerId, metaData.copy(), tbMsg.dataType, |
|
|
|
data, tbMsg.ruleChainId, tbMsg.ruleNodeId, tbMsg.ctx.copy(), tbMsg.callback); |
|
|
|
} |
|
|
|
|
|
|
|
public static TbMsg transformMsgOriginator(TbMsg tbMsg, EntityId originatorId) { |
|
|
|
return new TbMsg(tbMsg.queueName, tbMsg.id, tbMsg.ts, tbMsg.type, originatorId, tbMsg.getCustomerId(), tbMsg.metaData, tbMsg.dataType, |
|
|
|
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.ctx.copy(), tbMsg.getCallback()); |
|
|
|
} |
|
|
|
|
|
|
|
public static TbMsg transformMsgData(TbMsg tbMsg, String data) { |
|
|
|
return new TbMsg(tbMsg.queueName, tbMsg.id, tbMsg.ts, tbMsg.type, tbMsg.originator, tbMsg.customerId, tbMsg.metaData, tbMsg.dataType, |
|
|
|
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.ctx.copy(), tbMsg.getCallback()); |
|
|
|
} |
|
|
|
|
|
|
|
public static TbMsg transformMsgMetadata(TbMsg tbMsg, TbMsgMetaData metadata) { |
|
|
|
return new TbMsg(tbMsg.queueName, tbMsg.id, tbMsg.ts, tbMsg.type, tbMsg.originator, tbMsg.customerId, metadata.copy(), tbMsg.dataType, |
|
|
|
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.ctx.copy(), tbMsg.getCallback()); |
|
|
|
} |
|
|
|
|
|
|
|
public static TbMsg transformMsg(TbMsg tbMsg, TbMsgMetaData metadata, String data) { |
|
|
|
return new TbMsg(tbMsg.queueName, tbMsg.id, tbMsg.ts, tbMsg.type, tbMsg.originator, tbMsg.customerId, metadata, tbMsg.dataType, |
|
|
|
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.ctx.copy(), tbMsg.getCallback()); |
|
|
|
} |
|
|
|
|
|
|
|
public static TbMsg transformMsgCustomerId(TbMsg tbMsg, CustomerId customerId) { |
|
|
|
return new TbMsg(tbMsg.queueName, tbMsg.id, tbMsg.ts, tbMsg.type, tbMsg.originator, customerId, tbMsg.metaData, tbMsg.dataType, |
|
|
|
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.ctx.copy(), tbMsg.getCallback()); |
|
|
|
} |
|
|
|
|
|
|
|
public static TbMsg transformMsgRuleChainId(TbMsg tbMsg, RuleChainId ruleChainId) { |
|
|
|
return new TbMsg(tbMsg.queueName, tbMsg.id, tbMsg.ts, tbMsg.type, tbMsg.originator, tbMsg.customerId, tbMsg.metaData, tbMsg.dataType, |
|
|
|
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.ctx.copy(), tbMsg.getCallback()); |
|
|
|
} |
|
|
|
|
|
|
|
public static TbMsg transformMsgQueueName(TbMsg tbMsg, String queueName) { |
|
|
|
return new TbMsg(queueName, tbMsg.id, tbMsg.ts, tbMsg.type, tbMsg.originator, tbMsg.customerId, tbMsg.metaData, tbMsg.dataType, |
|
|
|
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.ctx.copy(), tbMsg.getCallback()); |
|
|
|
} |
|
|
|
|
|
|
|
public static TbMsg transformMsg(TbMsg tbMsg, RuleChainId ruleChainId, String queueName) { |
|
|
|
return new TbMsg(queueName, tbMsg.id, tbMsg.ts, tbMsg.type, tbMsg.originator, tbMsg.customerId, tbMsg.metaData, tbMsg.dataType, |
|
|
|
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.ctx.copy(), tbMsg.getCallback()); |
|
|
|
} |
|
|
|
|
|
|
|
//used for enqueueForTellNext
|
|
|
|
public static TbMsg newMsg(TbMsg tbMsg, String queueName, RuleChainId ruleChainId, RuleNodeId ruleNodeId) { |
|
|
|
return new TbMsg(queueName, UUID.randomUUID(), tbMsg.getTs(), tbMsg.getType(), tbMsg.getOriginator(), tbMsg.customerId, tbMsg.getMetaData().copy(), |
|
|
|
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.ctx.copy(), TbMsgCallback.EMPTY); |
|
|
|
} |
|
|
|
|
|
|
|
private TbMsg(String queueName, UUID id, long ts, String type, EntityId originator, CustomerId customerId, TbMsgMetaData metaData, TbMsgDataType dataType, String data, |
|
|
|
private TbMsg(String queueName, UUID id, long ts, TbMsgType internalType, EntityId originator, CustomerId customerId, TbMsgMetaData metaData, TbMsgDataType dataType, String data, |
|
|
|
RuleChainId ruleChainId, RuleNodeId ruleNodeId, TbMsgProcessingCtx ctx, TbMsgCallback callback) { |
|
|
|
this(queueName, id, ts, internalType, internalType.name(), originator, customerId, metaData, dataType, data, ruleChainId, ruleNodeId, ctx, callback); |
|
|
|
} |
|
|
|
|
|
|
|
private TbMsg(String queueName, UUID id, long ts, TbMsgType internalType, String type, EntityId originator, CustomerId customerId, TbMsgMetaData metaData, TbMsgDataType dataType, String data, |
|
|
|
RuleChainId ruleChainId, RuleNodeId ruleNodeId, TbMsgProcessingCtx ctx, TbMsgCallback callback) { |
|
|
|
this.id = id; |
|
|
|
this.queueName = queueName; |
|
|
|
@ -325,6 +331,7 @@ public final class TbMsg implements Serializable { |
|
|
|
this.ts = System.currentTimeMillis(); |
|
|
|
} |
|
|
|
this.type = type; |
|
|
|
this.internalType = internalType != null ? internalType : getInternalType(type); |
|
|
|
this.originator = originator; |
|
|
|
if (customerId == null || customerId.isNullUid()) { |
|
|
|
if (originator != null && originator.getEntityType() == EntityType.CUSTOMER) { |
|
|
|
@ -410,7 +417,7 @@ public final class TbMsg implements Serializable { |
|
|
|
} |
|
|
|
|
|
|
|
TbMsgDataType dataType = TbMsgDataType.values()[proto.getDataType()]; |
|
|
|
return new TbMsg(queueName, UUID.fromString(proto.getId()), proto.getTs(), proto.getType(), entityId, customerId, |
|
|
|
return new TbMsg(queueName, UUID.fromString(proto.getId()), proto.getTs(), null, proto.getType(), entityId, customerId, |
|
|
|
metaData, dataType, proto.getData(), ruleChainId, ruleNodeId, ctx, callback); |
|
|
|
} catch (InvalidProtocolBufferException e) { |
|
|
|
throw new IllegalStateException("Could not parse protobuf for TbMsg", e); |
|
|
|
@ -422,17 +429,17 @@ public final class TbMsg implements Serializable { |
|
|
|
} |
|
|
|
|
|
|
|
public TbMsg copyWithRuleChainId(RuleChainId ruleChainId, UUID msgId) { |
|
|
|
return new TbMsg(this.queueName, msgId, this.ts, this.type, this.originator, this.customerId, |
|
|
|
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.ctx, callback); |
|
|
|
} |
|
|
|
|
|
|
|
public TbMsg copyWithRuleNodeId(RuleChainId ruleChainId, RuleNodeId ruleNodeId, UUID msgId) { |
|
|
|
return new TbMsg(this.queueName, msgId, this.ts, this.type, this.originator, this.customerId, |
|
|
|
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.ctx, callback); |
|
|
|
} |
|
|
|
|
|
|
|
public TbMsg copyWithNewCtx() { |
|
|
|
return new TbMsg(this.queueName, this.id, this.ts, this.type, this.originator, this.customerId, |
|
|
|
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.ctx.copy(), TbMsgCallback.EMPTY); |
|
|
|
} |
|
|
|
|
|
|
|
@ -468,8 +475,16 @@ public final class TbMsg implements Serializable { |
|
|
|
return ts; |
|
|
|
} |
|
|
|
|
|
|
|
private TbMsgType getInternalType(String type) { |
|
|
|
try { |
|
|
|
return TbMsgType.valueOf(type); |
|
|
|
} catch (IllegalArgumentException e) { |
|
|
|
return TbMsgType.NA; |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
public boolean isTypeOf(TbMsgType tbMsgType) { |
|
|
|
return tbMsgType != null && tbMsgType.name().equals(this.type); |
|
|
|
return internalType.equals(tbMsgType); |
|
|
|
} |
|
|
|
|
|
|
|
public boolean isTypeOneOf(TbMsgType... types) { |
|
|
|
|