@ -15,43 +15,113 @@
* /
* /
package org.thingsboard.server.service.edge.rpc.constructor ;
package org.thingsboard.server.service.edge.rpc.constructor ;
import com.datastax.driver.core.utils.UUIDs ;
import com.fasterxml.jackson.core.JsonProcessingException ;
import com.fasterxml.jackson.core.JsonProcessingException ;
import com.fasterxml.jackson.databind.JsonNode ;
import com.fasterxml.jackson.databind.JsonNode ;
import com.fasterxml.jackson.databind.ObjectMapper ;
import lombok.extern.slf4j.Slf4j ;
import lombok.extern.slf4j.Slf4j ;
import org.jetbrains.annotations.NotNull ;
import org.jetbrains.annotations.NotNull ;
import org.junit.Assert ;
import org.junit.Assert ;
import org.junit.Before ;
import org.junit.Test ;
import org.junit.Test ;
import org.junit.runner.RunWith ;
import org.junit.runner.RunWith ;
import org.mockito.Mockito ;
import org.mockito.junit.MockitoJUnitRunner ;
import org.mockito.junit.MockitoJUnitRunner ;
import org.thingsboard.common.util.JacksonUtil ;
import org.thingsboard.server.common.data.id.QueueId ;
import org.thingsboard.server.common.data.id.RuleChainId ;
import org.thingsboard.server.common.data.id.RuleChainId ;
import org.thingsboard.server.common.data.id.RuleNodeId ;
import org.thingsboard.server.common.data.id.RuleNodeId ;
import org.thingsboard.server.common.data.id.TenantId ;
import org.thingsboard.server.common.data.queue.Queue ;
import org.thingsboard.server.common.data.rule.NodeConnectionInfo ;
import org.thingsboard.server.common.data.rule.NodeConnectionInfo ;
import org.thingsboard.server.common.data.rule.RuleChainMetaData ;
import org.thingsboard.server.common.data.rule.RuleChainMetaData ;
import org.thingsboard.server.common.data.rule.RuleNode ;
import org.thingsboard.server.common.data.rule.RuleNode ;
import org.thingsboard.server.dao.queue.QueueService ;
import org.thingsboard.server.gen.edge.v1.EdgeVersion ;
import org.thingsboard.server.gen.edge.v1.EdgeVersion ;
import org.thingsboard.server.gen.edge.v1.RuleChainConnectionInfoProto ;
import org.thingsboard.server.gen.edge.v1.RuleChainConnectionInfoProto ;
import org.thingsboard.server.gen.edge.v1.RuleChainMetadataUpdateMsg ;
import org.thingsboard.server.gen.edge.v1.RuleChainMetadataUpdateMsg ;
import org.thingsboard.server.gen.edge.v1.RuleNodeProto ;
import org.thingsboard.server.gen.edge.v1.UpdateMsgType ;
import org.thingsboard.server.gen.edge.v1.UpdateMsgType ;
import org.thingsboard.server.service.edge.rpc.constructor.rule.RuleChainMetadataConstructorV333 ;
import java.util.ArrayList ;
import java.util.ArrayList ;
import java.util.List ;
import java.util.List ;
import java.util.Optional ;
import java.util.UUID ;
import java.util.UUID ;
import static org.mockito.Mockito.mock ;
@Slf4j
@Slf4j
@RunWith ( MockitoJUnitRunner . class )
@RunWith ( MockitoJUnitRunner . class )
public class RuleChainMsgConstructorTest {
public class RuleChainMsgConstructorTest {
private static final ObjectMapper mapper = new ObjectMapper ( ) ;
private RuleChainMsgConstructor constructor ;
private QueueService queueService ;
private TenantId tenantId ;
private String queueId = "af588000-6c7c-11ec-bafd-c9a47a5c8d99" ;
@Before
public void setup ( ) {
queueService = mock ( QueueService . class ) ;
constructor = new RuleChainMsgConstructor ( queueService ) ;
tenantId = new TenantId ( UUID . randomUUID ( ) ) ;
}
@Test
public void testConstructRuleChainMetadataUpdatedMsg_V_3_4_0 ( ) throws JsonProcessingException {
RuleChainId ruleChainId = new RuleChainId ( UUID . randomUUID ( ) ) ;
RuleChainMetaData ruleChainMetaData = createRuleChainMetaData (
ruleChainId , 3 , createRuleNodes ( ruleChainId ) , createConnections ( ) ) ;
RuleChainMetadataUpdateMsg ruleChainMetadataUpdateMsg =
constructor . constructRuleChainMetadataUpdatedMsg (
tenantId ,
UpdateMsgType . ENTITY_CREATED_RPC_MESSAGE ,
ruleChainMetaData ,
EdgeVersion . V_3_4_0 ) ;
assetV_3_3_3_and_V_3_4_0 ( ruleChainMetadataUpdateMsg ) ;
assertCheckpointRuleNodeConfiguration (
ruleChainMetadataUpdateMsg . getNodesList ( ) ,
"{\"queueId\":\"" + queueId + "\"}" ) ;
}
@Test
@Test
public void testConstructRuleChainMetadataUpdatedMsg_V_3_3_3 ( ) throws JsonProcessingException {
public void testConstructRuleChainMetadataUpdatedMsg_V_3_3_3 ( ) throws JsonProcessingException {
Queue queue = new Queue ( ) ;
queue . setName ( "HighPriority" ) ;
Mockito . when ( queueService . findQueueById ( tenantId , new QueueId ( UUID . fromString ( queueId ) ) ) ) . thenReturn ( queue ) ;
RuleChainId ruleChainId = new RuleChainId ( UUID . randomUUID ( ) ) ;
RuleChainId ruleChainId = new RuleChainId ( UUID . randomUUID ( ) ) ;
RuleChainMsgConstructor constructor = new RuleChainMsgConstructor ( ) ;
RuleChainMetaData ruleChainMetaData = createRuleChainMetaData (
RuleChainMetaData ruleChainMetaData = createRuleChainMetaData ( ruleChainId , 3 , createRuleNodes ( ruleChainId ) , createConnections ( ) ) ;
ruleChainId , 3 , createRuleNodes ( ruleChainId ) , createConnections ( ) ) ;
RuleChainMetadataUpdateMsg ruleChainMetadataUpdateMsg =
RuleChainMetadataUpdateMsg ruleChainMetadataUpdateMsg =
constructor . constructRuleChainMetadataUpdatedMsg ( UpdateMsgType . ENTITY_CREATED_RPC_MESSAGE , ruleChainMetaData , EdgeVersion . V_3_3_3 ) ;
constructor . constructRuleChainMetadataUpdatedMsg (
tenantId ,
UpdateMsgType . ENTITY_CREATED_RPC_MESSAGE ,
ruleChainMetaData ,
EdgeVersion . V_3_3_3 ) ;
assetV_3_3_3_and_V_3_4_0 ( ruleChainMetadataUpdateMsg ) ;
assertCheckpointRuleNodeConfiguration (
ruleChainMetadataUpdateMsg . getNodesList ( ) ,
"{\"queueName\":\"HighPriority\"}" ) ;
}
private void assertCheckpointRuleNodeConfiguration ( List < RuleNodeProto > nodesList ,
String expectedConfiguration ) {
Optional < RuleNodeProto > checkpointRuleNodeOpt = nodesList . stream ( )
. filter ( rn - > RuleChainMetadataConstructorV333 . CHECKPOINT_NODE . equals ( rn . getType ( ) ) )
. findFirst ( ) ;
Assert . assertTrue ( checkpointRuleNodeOpt . isPresent ( ) ) ;
RuleNodeProto checkpointRuleNode = checkpointRuleNodeOpt . get ( ) ;
Assert . assertEquals ( expectedConfiguration , checkpointRuleNode . getConfiguration ( ) ) ;
}
private void assetV_3_3_3_and_V_3_4_0 ( RuleChainMetadataUpdateMsg ruleChainMetadataUpdateMsg ) {
Assert . assertEquals ( "First rule node index incorrect!" , 3 , ruleChainMetadataUpdateMsg . getFirstNodeIndex ( ) ) ;
Assert . assertEquals ( "First rule node index incorrect!" , 3 , ruleChainMetadataUpdateMsg . getFirstNodeIndex ( ) ) ;
Assert . assertEquals ( "Nodes count incorrect!" , 12 , ruleChainMetadataUpdateMsg . getNodesCount ( ) ) ;
Assert . assertEquals ( "Nodes count incorrect!" , 12 , ruleChainMetadataUpdateMsg . getNodesCount ( ) ) ;
Assert . assertEquals ( "Connections count incorrect!" , 13 , ruleChainMetadataUpdateMsg . getConnectionsCount ( ) ) ;
Assert . assertEquals ( "Connections count incorrect!" , 13 , ruleChainMetadataUpdateMsg . getConnectionsCount ( ) ) ;
@ -75,10 +145,13 @@ public class RuleChainMsgConstructorTest {
@Test
@Test
public void testConstructRuleChainMetadataUpdatedMsg_V_3_3_0 ( ) throws JsonProcessingException {
public void testConstructRuleChainMetadataUpdatedMsg_V_3_3_0 ( ) throws JsonProcessingException {
RuleChainId ruleChainId = new RuleChainId ( UUID . randomUUID ( ) ) ;
RuleChainId ruleChainId = new RuleChainId ( UUID . randomUUID ( ) ) ;
RuleChainMsgConstructor constructor = new RuleChainMsgConstructor ( ) ;
RuleChainMetaData ruleChainMetaData = createRuleChainMetaData ( ruleChainId , 3 , createRuleNodes ( ruleChainId ) , createConnections ( ) ) ;
RuleChainMetaData ruleChainMetaData = createRuleChainMetaData ( ruleChainId , 3 , createRuleNodes ( ruleChainId ) , createConnections ( ) ) ;
RuleChainMetadataUpdateMsg ruleChainMetadataUpdateMsg =
RuleChainMetadataUpdateMsg ruleChainMetadataUpdateMsg =
constructor . constructRuleChainMetadataUpdatedMsg ( UpdateMsgType . ENTITY_CREATED_RPC_MESSAGE , ruleChainMetaData , EdgeVersion . V_3_3_0 ) ;
constructor . constructRuleChainMetadataUpdatedMsg (
tenantId ,
UpdateMsgType . ENTITY_CREATED_RPC_MESSAGE ,
ruleChainMetaData ,
EdgeVersion . V_3_3_0 ) ;
Assert . assertEquals ( "First rule node index incorrect!" , 2 , ruleChainMetadataUpdateMsg . getFirstNodeIndex ( ) ) ;
Assert . assertEquals ( "First rule node index incorrect!" , 2 , ruleChainMetadataUpdateMsg . getFirstNodeIndex ( ) ) ;
Assert . assertEquals ( "Nodes count incorrect!" , 10 , ruleChainMetadataUpdateMsg . getNodesCount ( ) ) ;
Assert . assertEquals ( "Nodes count incorrect!" , 10 , ruleChainMetadataUpdateMsg . getNodesCount ( ) ) ;
@ -110,10 +183,13 @@ public class RuleChainMsgConstructorTest {
public void testConstructRuleChainMetadataUpdatedMsg_V_3_3_0_inDifferentOrder ( ) throws JsonProcessingException {
public void testConstructRuleChainMetadataUpdatedMsg_V_3_3_0_inDifferentOrder ( ) throws JsonProcessingException {
// same rule chain metadata, but different order of rule nodes
// same rule chain metadata, but different order of rule nodes
RuleChainId ruleChainId = new RuleChainId ( UUID . randomUUID ( ) ) ;
RuleChainId ruleChainId = new RuleChainId ( UUID . randomUUID ( ) ) ;
RuleChainMsgConstructor constructor = new RuleChainMsgConstructor ( ) ;
RuleChainMetaData ruleChainMetaData1 = createRuleChainMetaData ( ruleChainId , 8 , createRuleNodesInDifferentOrder ( ruleChainId ) , createConnectionsInDifferentOrder ( ) ) ;
RuleChainMetaData ruleChainMetaData1 = createRuleChainMetaData ( ruleChainId , 8 , createRuleNodesInDifferentOrder ( ruleChainId ) , createConnectionsInDifferentOrder ( ) ) ;
RuleChainMetadataUpdateMsg ruleChainMetadataUpdateMsg =
RuleChainMetadataUpdateMsg ruleChainMetadataUpdateMsg =
constructor . constructRuleChainMetadataUpdatedMsg ( UpdateMsgType . ENTITY_CREATED_RPC_MESSAGE , ruleChainMetaData1 , EdgeVersion . V_3_3_0 ) ;
constructor . constructRuleChainMetadataUpdatedMsg (
tenantId ,
UpdateMsgType . ENTITY_CREATED_RPC_MESSAGE ,
ruleChainMetaData1 ,
EdgeVersion . V_3_3_0 ) ;
Assert . assertEquals ( "First rule node index incorrect!" , 7 , ruleChainMetadataUpdateMsg . getFirstNodeIndex ( ) ) ;
Assert . assertEquals ( "First rule node index incorrect!" , 7 , ruleChainMetadataUpdateMsg . getFirstNodeIndex ( ) ) ;
Assert . assertEquals ( "Nodes count incorrect!" , 10 , ruleChainMetadataUpdateMsg . getNodesCount ( ) ) ;
Assert . assertEquals ( "Nodes count incorrect!" , 10 , ruleChainMetadataUpdateMsg . getNodesCount ( ) ) ;
@ -253,8 +329,8 @@ public class RuleChainMsgConstructorTest {
return createRuleNode ( ruleChainId ,
return createRuleNode ( ruleChainId ,
"org.thingsboard.rule.engine.flow.TbRuleChainOutputNode" ,
"org.thingsboard.rule.engine.flow.TbRuleChainOutputNode" ,
"Output node" ,
"Output node" ,
mapper . readTree ( "{\"version\":0}" ) ,
JacksonUtil . OBJECT_MAPPER . readTree ( "{\"version\":0}" ) ,
mapper . readTree ( "{\"description\":\"\",\"layoutX\":178,\"layoutY\":592}" ) ) ;
JacksonUtil . OBJECT_MAPPER . readTree ( "{\"description\":\"\",\"layoutX\":178,\"layoutY\":592}" ) ) ;
}
}
@NotNull
@NotNull
@ -262,8 +338,8 @@ public class RuleChainMsgConstructorTest {
return createRuleNode ( ruleChainId ,
return createRuleNode ( ruleChainId ,
"org.thingsboard.rule.engine.flow.TbCheckpointNode" ,
"org.thingsboard.rule.engine.flow.TbCheckpointNode" ,
"Checkpoint node" ,
"Checkpoint node" ,
mapper . readTree ( "{\"queueName\":\"HighPriority \"}" ) ,
JacksonUtil . OBJECT_MAPPER . readTree ( "{\"queueId\":\"" + queueId + " \"}" ) ,
mapper . readTree ( "{\"description\":\"\",\"layoutX\":178,\"layoutY\":647}" ) ) ;
JacksonUtil . OBJECT_MAPPER . readTree ( "{\"description\":\"\",\"layoutX\":178,\"layoutY\":647}" ) ) ;
}
}
@NotNull
@NotNull
@ -271,8 +347,8 @@ public class RuleChainMsgConstructorTest {
return createRuleNode ( ruleChainId ,
return createRuleNode ( ruleChainId ,
"org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNode" ,
"org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNode" ,
"Save Timeseries" ,
"Save Timeseries" ,
mapper . readTree ( "{\"defaultTTL\":0}" ) ,
JacksonUtil . OBJECT_MAPPER . readTree ( "{\"defaultTTL\":0}" ) ,
mapper . readTree ( "{\"layoutX\":823,\"layoutY\":157}" ) ) ;
JacksonUtil . OBJECT_MAPPER . readTree ( "{\"layoutX\":823,\"layoutY\":157}" ) ) ;
}
}
@NotNull
@NotNull
@ -280,8 +356,8 @@ public class RuleChainMsgConstructorTest {
return createRuleNode ( ruleChainId ,
return createRuleNode ( ruleChainId ,
"org.thingsboard.rule.engine.filter.TbMsgTypeSwitchNode" ,
"org.thingsboard.rule.engine.filter.TbMsgTypeSwitchNode" ,
"Message Type Switch" ,
"Message Type Switch" ,
mapper . readTree ( "{\"version\":0}" ) ,
JacksonUtil . OBJECT_MAPPER . readTree ( "{\"version\":0}" ) ,
mapper . readTree ( "{\"layoutX\":347,\"layoutY\":149}" ) ) ;
JacksonUtil . OBJECT_MAPPER . readTree ( "{\"layoutX\":347,\"layoutY\":149}" ) ) ;
}
}
@NotNull
@NotNull
@ -289,8 +365,8 @@ public class RuleChainMsgConstructorTest {
return createRuleNode ( ruleChainId ,
return createRuleNode ( ruleChainId ,
"org.thingsboard.rule.engine.action.TbLogNode" ,
"org.thingsboard.rule.engine.action.TbLogNode" ,
"Log Other" ,
"Log Other" ,
mapper . readTree ( "{\"jsScript\":\"return '\\\\nIncoming message:\\\\n' + JSON.stringify(msg) + '\\\\nIncoming metadata:\\\\n' + JSON.stringify(metadata);\"}" ) ,
JacksonUtil . OBJECT_MAPPER . readTree ( "{\"jsScript\":\"return '\\\\nIncoming message:\\\\n' + JSON.stringify(msg) + '\\\\nIncoming metadata:\\\\n' + JSON.stringify(metadata);\"}" ) ,
mapper . readTree ( "{\"layoutX\":824,\"layoutY\":378}" ) ) ;
JacksonUtil . OBJECT_MAPPER . readTree ( "{\"layoutX\":824,\"layoutY\":378}" ) ) ;
}
}
@NotNull
@NotNull
@ -298,8 +374,8 @@ public class RuleChainMsgConstructorTest {
return createRuleNode ( ruleChainId ,
return createRuleNode ( ruleChainId ,
"org.thingsboard.rule.engine.edge.TbMsgPushToCloudNode" ,
"org.thingsboard.rule.engine.edge.TbMsgPushToCloudNode" ,
"Push to cloud" ,
"Push to cloud" ,
mapper . readTree ( "{\"scope\":\"SERVER_SCOPE\"}" ) ,
JacksonUtil . OBJECT_MAPPER . readTree ( "{\"scope\":\"SERVER_SCOPE\"}" ) ,
mapper . readTree ( "{\"layoutX\":1129,\"layoutY\":52}" ) ) ;
JacksonUtil . OBJECT_MAPPER . readTree ( "{\"layoutX\":1129,\"layoutY\":52}" ) ) ;
}
}
@NotNull
@NotNull
@ -307,8 +383,8 @@ public class RuleChainMsgConstructorTest {
return createRuleNode ( ruleChainId ,
return createRuleNode ( ruleChainId ,
"org.thingsboard.rule.engine.flow.TbAckNode" ,
"org.thingsboard.rule.engine.flow.TbAckNode" ,
"Acknowledge node" ,
"Acknowledge node" ,
mapper . readTree ( "{\"version\":0}" ) ,
JacksonUtil . OBJECT_MAPPER . readTree ( "{\"version\":0}" ) ,
mapper . readTree ( "{\"description\":\"\",\"layoutX\":177,\"layoutY\":703}" ) ) ;
JacksonUtil . OBJECT_MAPPER . readTree ( "{\"description\":\"\",\"layoutX\":177,\"layoutY\":703}" ) ) ;
}
}
@NotNull
@NotNull
@ -316,8 +392,8 @@ public class RuleChainMsgConstructorTest {
return createRuleNode ( ruleChainId ,
return createRuleNode ( ruleChainId ,
"org.thingsboard.rule.engine.profile.TbDeviceProfileNode" ,
"org.thingsboard.rule.engine.profile.TbDeviceProfileNode" ,
"Device Profile Node" ,
"Device Profile Node" ,
mapper . readTree ( "{\"persistAlarmRulesState\":false,\"fetchAlarmRulesStateOnStart\":false}" ) ,
JacksonUtil . OBJECT_MAPPER . readTree ( "{\"persistAlarmRulesState\":false,\"fetchAlarmRulesStateOnStart\":false}" ) ,
mapper . readTree ( "{\"description\":\"Process incoming messages from devices with the alarm rules defined in the device profile. Dispatch all incoming messages with \\\"Success\\\" relation type.\",\"layoutX\":187,\"layoutY\":468}" ) ) ;
JacksonUtil . OBJECT_MAPPER . readTree ( "{\"description\":\"Process incoming messages from devices with the alarm rules defined in the device profile. Dispatch all incoming messages with \\\"Success\\\" relation type.\",\"layoutX\":187,\"layoutY\":468}" ) ) ;
}
}
@NotNull
@NotNull
@ -325,8 +401,8 @@ public class RuleChainMsgConstructorTest {
return createRuleNode ( ruleChainId ,
return createRuleNode ( ruleChainId ,
"org.thingsboard.rule.engine.telemetry.TbMsgAttributesNode" ,
"org.thingsboard.rule.engine.telemetry.TbMsgAttributesNode" ,
"Save Client Attributes" ,
"Save Client Attributes" ,
mapper . readTree ( "{\"scope\":\"CLIENT_SCOPE\"}" ) ,
JacksonUtil . OBJECT_MAPPER . readTree ( "{\"scope\":\"CLIENT_SCOPE\"}" ) ,
mapper . readTree ( "{\"layoutX\":824,\"layoutY\":52}" ) ) ;
JacksonUtil . OBJECT_MAPPER . readTree ( "{\"layoutX\":824,\"layoutY\":52}" ) ) ;
}
}
@NotNull
@NotNull
@ -334,8 +410,8 @@ public class RuleChainMsgConstructorTest {
return createRuleNode ( ruleChainId ,
return createRuleNode ( ruleChainId ,
"org.thingsboard.rule.engine.action.TbLogNode" ,
"org.thingsboard.rule.engine.action.TbLogNode" ,
"Log RPC from Device" ,
"Log RPC from Device" ,
mapper . readTree ( "{\"jsScript\":\"return '\\\\nIncoming message:\\\\n' + JSON.stringify(msg) + '\\\\nIncoming metadata:\\\\n' + JSON.stringify(metadata);\"}" ) ,
JacksonUtil . OBJECT_MAPPER . readTree ( "{\"jsScript\":\"return '\\\\nIncoming message:\\\\n' + JSON.stringify(msg) + '\\\\nIncoming metadata:\\\\n' + JSON.stringify(metadata);\"}" ) ,
mapper . readTree ( "{\"layoutX\":825,\"layoutY\":266}" ) ) ;
JacksonUtil . OBJECT_MAPPER . readTree ( "{\"layoutX\":825,\"layoutY\":266}" ) ) ;
}
}
@NotNull
@NotNull
@ -343,8 +419,8 @@ public class RuleChainMsgConstructorTest {
return createRuleNode ( ruleChainId ,
return createRuleNode ( ruleChainId ,
"org.thingsboard.rule.engine.rpc.TbSendRPCRequestNode" ,
"org.thingsboard.rule.engine.rpc.TbSendRPCRequestNode" ,
"RPC Call Request" ,
"RPC Call Request" ,
mapper . readTree ( "{\"timeoutInSeconds\":60}" ) ,
JacksonUtil . OBJECT_MAPPER . readTree ( "{\"timeoutInSeconds\":60}" ) ,
mapper . readTree ( "{\"layoutX\":824,\"layoutY\":466}" ) ) ;
JacksonUtil . OBJECT_MAPPER . readTree ( "{\"layoutX\":824,\"layoutY\":466}" ) ) ;
}
}
@NotNull
@NotNull
@ -352,7 +428,7 @@ public class RuleChainMsgConstructorTest {
return createRuleNode ( ruleChainId ,
return createRuleNode ( ruleChainId ,
"org.thingsboard.rule.engine.flow.TbRuleChainInputNode" ,
"org.thingsboard.rule.engine.flow.TbRuleChainInputNode" ,
"Push to Analytics" ,
"Push to Analytics" ,
mapper . readTree ( "{\"ruleChainId\":\"af588000-6c7c-11ec-bafd-c9a47a5c8d99\"}" ) ,
JacksonUtil . OBJECT_MAPPER . readTree ( "{\"ruleChainId\":\"af588000-6c7c-11ec-bafd-c9a47a5c8d99\"}" ) ,
mapper . readTree ( "{\"description\":\"\",\"layoutX\":477,\"layoutY\":560}" ) ) ;
JacksonUtil . OBJECT_MAPPER . readTree ( "{\"description\":\"\",\"layoutX\":477,\"layoutY\":560}" ) ) ;
}
}
}
}