@ -18,27 +18,27 @@ package org.thingsboard.server.common.msg;
import com.fasterxml.jackson.annotation.JsonIgnore ;
import com.google.protobuf.ByteString ;
import com.google.protobuf.InvalidProtocolBufferException ;
import lombok.AccessLevel ;
import lombok.Builder ;
import lombok.Data ;
import lombok.Getter ;
import lombok.extern.slf4j.Slf4j ;
import org.thingsboard.server.common.data.id.EntityId ;
import org.thingsboard.server.common.data.id.EntityIdFactory ;
import org.thingsboard.server.common.data.id.RuleChainId ;
import org.thingsboard.server.common.data.id.RuleNodeId ;
import org.thingsboard.server.common.msg.gen.MsgProtos ;
import org.thingsboard.server.common.msg.queue.RuleNodeInfo ;
import org.thingsboard.server.common.msg.queue.ServiceQueue ;
import org.thingsboard.server.common.msg.queue.TbMsgCallback ;
import java.io.IOException ;
import java.io.Serializable ;
import java.util.UUID ;
import java.util.concurrent.atomic.AtomicInteger ;
/ * *
* Created by ashvayka on 13 . 01 . 18 .
* /
@Data
@Builder
@Slf4j
public final class TbMsg implements Serializable {
@ -52,51 +52,63 @@ public final class TbMsg implements Serializable {
private final String data ;
private final RuleChainId ruleChainId ;
private final RuleNodeId ruleNodeId ;
@Getter ( value = AccessLevel . NONE )
private final AtomicInteger ruleNodeExecCounter ;
public int getAndIncrementRuleNodeCounter ( ) {
return ruleNodeExecCounter . getAndIncrement ( ) ;
}
//This field is not serialized because we use queues and there is no need to do it
@JsonIgnore
transient private final TbMsgCallback callback ;
public static TbMsg newMsg ( String queueName , String type , EntityId originator , TbMsgMetaData metaData , String data , RuleChainId ruleChainId , RuleNodeId ruleNodeId ) {
return new TbMsg ( queueName , UUID . randomUUID ( ) , System . currentTimeMillis ( ) , type , originator , metaData . copy ( ) , TbMsgDataType . JSON , data , ruleChainId , ruleNodeId , TbMsgCallback . EMPTY ) ;
return new TbMsg ( queueName , UUID . randomUUID ( ) , System . currentTimeMillis ( ) , type , originator ,
metaData . copy ( ) , TbMsgDataType . JSON , data , ruleChainId , ruleNodeId , 0 , TbMsgCallback . EMPTY ) ;
}
public static TbMsg newMsg ( String type , EntityId originator , TbMsgMetaData metaData , String data ) {
return new TbMsg ( ServiceQueue . MAIN , UUID . randomUUID ( ) , System . currentTimeMillis ( ) , type , originator , metaData . copy ( ) , TbMsgDataType . JSON , data , null , null , TbMsgCallback . EMPTY ) ;
return new TbMsg ( ServiceQueue . MAIN , UUID . randomUUID ( ) , System . currentTimeMillis ( ) , type , originator , metaData . copy ( ) , TbMsgDataType . JSON , data , null , null , 0 , TbMsgCallback . EMPTY ) ;
}
// REALLY NEW MSG
public static TbMsg newMsg ( String queueName , String type , EntityId originator , TbMsgMetaData metaData , String data ) {
return new TbMsg ( queueName , UUID . randomUUID ( ) , System . currentTimeMillis ( ) , type , originator , metaData . copy ( ) , TbMsgDataType . JSON , data , null , null , TbMsgCallback . EMPTY ) ;
return new TbMsg ( queueName , UUID . randomUUID ( ) , System . currentTimeMillis ( ) , type , originator , metaData . copy ( ) , TbMsgDataType . JSON , data , null , null , 0 , TbMsgCallback . EMPTY ) ;
}
public static TbMsg newMsg ( String type , EntityId originator , TbMsgMetaData metaData , TbMsgDataType dataType , String data ) {
return new TbMsg ( ServiceQueue . MAIN , UUID . randomUUID ( ) , System . currentTimeMillis ( ) , type , originator , metaData . copy ( ) , dataType , data , null , null , TbMsgCallback . EMPTY ) ;
return new TbMsg ( ServiceQueue . MAIN , UUID . randomUUID ( ) , System . currentTimeMillis ( ) , type , originator , metaData . copy ( ) , dataType , data , null , null , 0 , TbMsgCallback . EMPTY ) ;
}
// For Tests only
public static TbMsg newMsg ( String type , EntityId originator , TbMsgMetaData metaData , TbMsgDataType dataType , String data , RuleChainId ruleChainId , RuleNodeId ruleNodeId ) {
return new TbMsg ( ServiceQueue . MAIN , UUID . randomUUID ( ) , System . currentTimeMillis ( ) , type , originator , metaData . copy ( ) , dataType , data , ruleChainId , ruleNodeId , TbMsgCallback . EMPTY ) ;
return new TbMsg ( ServiceQueue . MAIN , UUID . randomUUID ( ) , System . currentTimeMillis ( ) , type , originator , metaData . copy ( ) , dataType , data , ruleChainId , ruleNodeId , 0 , TbMsgCallback . EMPTY ) ;
}
public static TbMsg newMsg ( String type , EntityId originator , TbMsgMetaData metaData , String data , TbMsgCallback callback ) {
return new TbMsg ( ServiceQueue . MAIN , UUID . randomUUID ( ) , System . currentTimeMillis ( ) , type , originator , metaData . copy ( ) , TbMsgDataType . JSON , data , null , null , callback ) ;
return new TbMsg ( ServiceQueue . MAIN , UUID . randomUUID ( ) , System . currentTimeMillis ( ) , type , originator , metaData . copy ( ) , TbMsgDataType . JSON , data , null , null , 0 , callback ) ;
}
public static TbMsg transformMsg ( TbMsg orig Msg, String type , EntityId originator , TbMsgMetaData metaData , String data ) {
return new TbMsg ( orig Msg. getQueueName ( ) , orig Msg. getId ( ) , orig Msg. getTs ( ) , type , originator , metaData . copy ( ) , orig Msg. getDataType ( ) ,
data , orig Msg. getRuleChainId ( ) , orig Msg. getRuleNodeId ( ) , orig Msg. getCallback ( ) ) ;
public static TbMsg transformMsg ( TbMsg tb Msg, String type , EntityId originator , TbMsgMetaData metaData , String data ) {
return new TbMsg ( tb Msg. getQueueName ( ) , tb Msg. getId ( ) , tb Msg. getTs ( ) , type , originator , metaData . copy ( ) , tb Msg. getDataType ( ) ,
data , tb Msg. getRuleChainId ( ) , tb Msg. getRuleNodeId ( ) , tbMsg . ruleNodeExecCounter . get ( ) , tb Msg. getCallback ( ) ) ;
}
public static TbMsg transformMsg ( TbMsg orig Msg, RuleChainId ruleChainId ) {
return new TbMsg ( orig Msg. queueName , orig Msg. id , orig Msg. ts , orig Msg. type , orig Msg. originator , orig Msg. metaData , orig Msg. dataType ,
orig Msg. data , ruleChainId , null , orig Msg. getCallback ( ) ) ;
public static TbMsg transformMsg ( TbMsg tb Msg, RuleChainId ruleChainId ) {
return new TbMsg ( tb Msg. queueName , tb Msg. id , tb Msg. ts , tb Msg. type , tb Msg. originator , tb Msg. metaData , tb Msg. dataType ,
tb Msg. data , ruleChainId , null , tbMsg . ruleNodeExecCounter . get ( ) , tb Msg. getCallback ( ) ) ;
}
public static TbMsg newMsg ( TbMsg tbMsg , RuleChainId ruleChainId , RuleNodeId ruleNodeId ) {
return new TbMsg ( tbMsg . getQueueName ( ) , UUID . randomUUID ( ) , tbMsg . getTs ( ) , tbMsg . getType ( ) , tbMsg . getOriginator ( ) , tbMsg . getMetaData ( ) . copy ( ) ,
tbMsg . getDataType ( ) , tbMsg . getData ( ) , ruleChainId , ruleNodeId , TbMsgCallback . EMPTY ) ;
tbMsg . getDataType ( ) , tbMsg . getData ( ) , ruleChainId , ruleNodeId , tbMsg . ruleNodeExecCounter . get ( ) , TbMsgCallback . EMPTY ) ;
}
private TbMsg ( String queueName , UUID id , long ts , String type , EntityId originator , TbMsgMetaData metaData , TbMsgDataType dataType , String data ,
RuleChainId ruleChainId , RuleNodeId ruleNodeId , TbMsgCallback callback ) {
RuleChainId ruleChainId , RuleNodeId ruleNodeId , int ruleNodeExecCounter , TbMsgCallback callback ) {
this . id = id ;
this . queueName = queueName ;
if ( ts > 0 ) {
@ -111,6 +123,7 @@ public final class TbMsg implements Serializable {
this . data = data ;
this . ruleChainId = ruleChainId ;
this . ruleNodeId = ruleNodeId ;
this . ruleNodeExecCounter = new AtomicInteger ( ruleNodeExecCounter ) ;
if ( callback ! = null ) {
this . callback = callback ;
} else {
@ -147,6 +160,7 @@ public final class TbMsg implements Serializable {
builder . setDataType ( msg . getDataType ( ) . ordinal ( ) ) ;
builder . setData ( msg . getData ( ) ) ;
builder . setRuleNodeExecCounter ( msg . ruleNodeExecCounter . get ( ) ) ;
return builder . build ( ) . toByteArray ( ) ;
}
@ -164,18 +178,18 @@ public final class TbMsg implements Serializable {
ruleNodeId = new RuleNodeId ( new UUID ( proto . getRuleNodeIdMSB ( ) , proto . getRuleNodeIdLSB ( ) ) ) ;
}
TbMsgDataType dataType = TbMsgDataType . values ( ) [ proto . getDataType ( ) ] ;
return new TbMsg ( queueName , UUID . fromString ( proto . getId ( ) ) , proto . getTs ( ) , proto . getType ( ) , entityId , metaData , dataType , proto . getData ( ) , ruleChainId , ruleNodeId , callback ) ;
return new TbMsg ( queueName , UUID . fromString ( proto . getId ( ) ) , proto . getTs ( ) , proto . getType ( ) , entityId , metaData , dataType , proto . getData ( ) , ruleChainId , ruleNodeId , proto . getRuleNodeExecCounter ( ) , callback ) ;
} catch ( InvalidProtocolBufferException e ) {
throw new IllegalStateException ( "Could not parse protobuf for TbMsg" , e ) ;
}
}
public TbMsg copyWithRuleChainId ( RuleChainId ruleChainId ) {
return new TbMsg ( this . queueName , this . id , this . ts , this . type , this . originator , this . metaData , this . dataType , this . data , ruleChainId , null , callback ) ;
return new TbMsg ( this . queueName , this . id , this . ts , this . type , this . originator , this . metaData , this . dataType , this . data , ruleChainId , null , this . ruleNodeExecCounter . get ( ) , callback ) ;
}
public TbMsg copyWithRuleNodeId ( RuleChainId ruleChainId , RuleNodeId ruleNodeId ) {
return new TbMsg ( this . queueName , this . id , this . ts , this . type , this . originator , this . metaData , this . dataType , this . data , ruleChainId , ruleNodeId , callback ) ;
return new TbMsg ( this . queueName , this . id , this . ts , this . type , this . originator , this . metaData , this . dataType , this . data , ruleChainId , ruleNodeId , this . ruleNodeExecCounter . get ( ) , callback ) ;
}
public TbMsgCallback getCallback ( ) {