|
|
@ -41,6 +41,7 @@ public final class TbMsg implements Serializable { |
|
|
private final TbMsgMetaData metaData; |
|
|
private final TbMsgMetaData metaData; |
|
|
private final TbMsgDataType dataType; |
|
|
private final TbMsgDataType dataType; |
|
|
private final String data; |
|
|
private final String data; |
|
|
|
|
|
private final TbMsgTransactionData transactionData; |
|
|
|
|
|
|
|
|
//The following fields are not persisted to DB, because they can always be recovered from the context;
|
|
|
//The following fields are not persisted to DB, because they can always be recovered from the context;
|
|
|
private final RuleChainId ruleChainId; |
|
|
private final RuleChainId ruleChainId; |
|
|
@ -55,11 +56,17 @@ public final class TbMsg implements Serializable { |
|
|
this.metaData = metaData; |
|
|
this.metaData = metaData; |
|
|
this.data = data; |
|
|
this.data = data; |
|
|
this.dataType = TbMsgDataType.JSON; |
|
|
this.dataType = TbMsgDataType.JSON; |
|
|
|
|
|
this.transactionData = null; |
|
|
this.ruleChainId = ruleChainId; |
|
|
this.ruleChainId = ruleChainId; |
|
|
this.ruleNodeId = ruleNodeId; |
|
|
this.ruleNodeId = ruleNodeId; |
|
|
this.clusterPartition = clusterPartition; |
|
|
this.clusterPartition = clusterPartition; |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
public TbMsg(UUID id, String type, EntityId originator, TbMsgMetaData metaData, TbMsgDataType dataType, String data, |
|
|
|
|
|
RuleChainId ruleChainId, RuleNodeId ruleNodeId, long clusterPartition) { |
|
|
|
|
|
this(id, type, originator, metaData, dataType, data, null, ruleChainId, ruleNodeId, clusterPartition); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
public static ByteBuffer toBytes(TbMsg msg) { |
|
|
public static ByteBuffer toBytes(TbMsg msg) { |
|
|
MsgProtos.TbMsgProto.Builder builder = MsgProtos.TbMsgProto.newBuilder(); |
|
|
MsgProtos.TbMsgProto.Builder builder = MsgProtos.TbMsgProto.newBuilder(); |
|
|
builder.setId(msg.getId().toString()); |
|
|
builder.setId(msg.getId().toString()); |
|
|
@ -82,6 +89,16 @@ public final class TbMsg implements Serializable { |
|
|
builder.setMetaData(MsgProtos.TbMsgMetaDataProto.newBuilder().putAllData(msg.getMetaData().getData()).build()); |
|
|
builder.setMetaData(MsgProtos.TbMsgMetaDataProto.newBuilder().putAllData(msg.getMetaData().getData()).build()); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
TbMsgTransactionData transactionData = msg.getTransactionData(); |
|
|
|
|
|
if (transactionData != null) { |
|
|
|
|
|
MsgProtos.TbMsgTransactionDataProto.Builder transactionBuilder = MsgProtos.TbMsgTransactionDataProto.newBuilder(); |
|
|
|
|
|
transactionBuilder.setId(transactionData.getTransactionId().toString()); |
|
|
|
|
|
transactionBuilder.setEntityType(transactionData.getOriginatorId().getEntityType().name()); |
|
|
|
|
|
transactionBuilder.setEntityIdMSB(transactionData.getOriginatorId().getId().getMostSignificantBits()); |
|
|
|
|
|
transactionBuilder.setEntityIdLSB(transactionData.getOriginatorId().getId().getLeastSignificantBits()); |
|
|
|
|
|
builder.setTransactionData(transactionBuilder.build()); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
builder.setDataType(msg.getDataType().ordinal()); |
|
|
builder.setDataType(msg.getDataType().ordinal()); |
|
|
builder.setData(msg.getData()); |
|
|
builder.setData(msg.getData()); |
|
|
byte[] bytes = builder.build().toByteArray(); |
|
|
byte[] bytes = builder.build().toByteArray(); |
|
|
@ -92,6 +109,9 @@ public final class TbMsg implements Serializable { |
|
|
try { |
|
|
try { |
|
|
MsgProtos.TbMsgProto proto = MsgProtos.TbMsgProto.parseFrom(buffer.array()); |
|
|
MsgProtos.TbMsgProto proto = MsgProtos.TbMsgProto.parseFrom(buffer.array()); |
|
|
TbMsgMetaData metaData = new TbMsgMetaData(proto.getMetaData().getDataMap()); |
|
|
TbMsgMetaData metaData = new TbMsgMetaData(proto.getMetaData().getDataMap()); |
|
|
|
|
|
EntityId transactionEntityId = EntityIdFactory.getByTypeAndUuid(proto.getTransactionData().getEntityType(), |
|
|
|
|
|
new UUID(proto.getTransactionData().getEntityIdMSB(), proto.getTransactionData().getEntityIdLSB())); |
|
|
|
|
|
TbMsgTransactionData transactionData = new TbMsgTransactionData(UUID.fromString(proto.getTransactionData().getId()), transactionEntityId); |
|
|
EntityId entityId = EntityIdFactory.getByTypeAndUuid(proto.getEntityType(), new UUID(proto.getEntityIdMSB(), proto.getEntityIdLSB())); |
|
|
EntityId entityId = EntityIdFactory.getByTypeAndUuid(proto.getEntityType(), new UUID(proto.getEntityIdMSB(), proto.getEntityIdLSB())); |
|
|
RuleChainId ruleChainId = new RuleChainId(new UUID(proto.getRuleChainIdMSB(), proto.getRuleChainIdLSB())); |
|
|
RuleChainId ruleChainId = new RuleChainId(new UUID(proto.getRuleChainIdMSB(), proto.getRuleChainIdLSB())); |
|
|
RuleNodeId ruleNodeId = null; |
|
|
RuleNodeId ruleNodeId = null; |
|
|
@ -99,7 +119,7 @@ public final class TbMsg implements Serializable { |
|
|
ruleNodeId = new RuleNodeId(new UUID(proto.getRuleNodeIdMSB(), proto.getRuleNodeIdLSB())); |
|
|
ruleNodeId = new RuleNodeId(new UUID(proto.getRuleNodeIdMSB(), proto.getRuleNodeIdLSB())); |
|
|
} |
|
|
} |
|
|
TbMsgDataType dataType = TbMsgDataType.values()[proto.getDataType()]; |
|
|
TbMsgDataType dataType = TbMsgDataType.values()[proto.getDataType()]; |
|
|
return new TbMsg(UUID.fromString(proto.getId()), proto.getType(), entityId, metaData, dataType, proto.getData(), ruleChainId, ruleNodeId, proto.getClusterPartition()); |
|
|
return new TbMsg(UUID.fromString(proto.getId()), proto.getType(), entityId, metaData, dataType, proto.getData(), transactionData, ruleChainId, ruleNodeId, proto.getClusterPartition()); |
|
|
} catch (InvalidProtocolBufferException e) { |
|
|
} catch (InvalidProtocolBufferException e) { |
|
|
throw new IllegalStateException("Could not parse protobuf for TbMsg", e); |
|
|
throw new IllegalStateException("Could not parse protobuf for TbMsg", e); |
|
|
} |
|
|
} |
|
|
|