@ -16,13 +16,9 @@
package org.thingsboard.rule.engine.rest ;
package org.thingsboard.rule.engine.rest ;
import com.datastax.oss.driver.api.core.uuid.Uuids ;
import com.datastax.oss.driver.api.core.uuid.Uuids ;
import org.apache.http.HttpException ;
import org.apache.http.HttpRequest ;
import org.apache.http.HttpResponse ;
import org.apache.http.config.SocketConfig ;
import org.apache.http.config.SocketConfig ;
import org.apache.http.impl.bootstrap.HttpServer ;
import org.apache.http.impl.bootstrap.HttpServer ;
import org.apache.http.impl.bootstrap.ServerBootstrap ;
import org.apache.http.impl.bootstrap.ServerBootstrap ;
import org.apache.http.protocol.HttpContext ;
import org.apache.http.protocol.HttpRequestHandler ;
import org.apache.http.protocol.HttpRequestHandler ;
import org.junit.jupiter.api.AfterEach ;
import org.junit.jupiter.api.AfterEach ;
import org.junit.jupiter.api.BeforeEach ;
import org.junit.jupiter.api.BeforeEach ;
@ -52,11 +48,14 @@ import org.thingsboard.server.common.msg.TbMsgMetaData;
import org.thingsboard.server.exception.DataValidationException ;
import org.thingsboard.server.exception.DataValidationException ;
import java.io.IOException ;
import java.io.IOException ;
import java.io.InputStream ;
import java.nio.charset.StandardCharsets ;
import java.util.ArrayList ;
import java.util.ArrayList ;
import java.util.Collections ;
import java.util.Collections ;
import java.util.List ;
import java.util.List ;
import java.util.concurrent.CountDownLatch ;
import java.util.concurrent.CountDownLatch ;
import java.util.concurrent.TimeUnit ;
import java.util.concurrent.TimeUnit ;
import java.util.concurrent.atomic.AtomicReference ;
import java.util.stream.Stream ;
import java.util.stream.Stream ;
import static org.assertj.core.api.Assertions.assertThatThrownBy ;
import static org.assertj.core.api.Assertions.assertThatThrownBy ;
@ -71,19 +70,17 @@ import static org.mockito.Mockito.verify;
@ExtendWith ( MockitoExtension . class )
@ExtendWith ( MockitoExtension . class )
public class TbRestApiCallNodeTest extends AbstractRuleNodeUpgradeTest {
public class TbRestApiCallNodeTest extends AbstractRuleNodeUpgradeTest {
private RuleNode ruleNode ;
@Spy
@Spy
private TbRestApiCallNode restNode ;
private TbRestApiCallNode restNode ;
@Mock
@Mock
private TbContext ctx ;
private TbContext ctx ;
private EntityId originator = new DeviceId ( Uuids . timeBased ( ) ) ;
private final EntityId originator = new DeviceId ( Uuids . timeBased ( ) ) ;
private TbMsgMetaData metaData = new TbMsgMetaData ( ) ;
private final TbMsgMetaData metaData = new TbMsgMetaData ( ) ;
private RuleChainId ruleChainId = new RuleChainId ( Uuids . timeBased ( ) ) ;
private final RuleChainId ruleChainId = new RuleChainId ( Uuids . timeBased ( ) ) ;
private RuleNodeId ruleNodeId = new RuleNodeId ( Uuids . timeBased ( ) ) ;
private final RuleNodeId ruleNodeId = new RuleNodeId ( Uuids . timeBased ( ) ) ;
private HttpServer server ;
private HttpServer server ;
@ -108,7 +105,7 @@ public class TbRestApiCallNodeTest extends AbstractRuleNodeUpgradeTest {
@BeforeEach
@BeforeEach
public void setup ( ) {
public void setup ( ) {
ruleNode = new RuleNode ( ) ;
RuleNode ruleNode = new RuleNode ( ) ;
ruleNode . setId ( ruleNodeId ) ;
ruleNode . setId ( ruleNodeId ) ;
ruleNode . setName ( "Test REST API call node" ) ;
ruleNode . setName ( "Test REST API call node" ) ;
lenient ( ) . when ( ctx . getSelf ( ) ) . thenReturn ( ruleNode ) ;
lenient ( ) . when ( ctx . getSelf ( ) ) . thenReturn ( ruleNode ) ;
@ -217,22 +214,17 @@ public class TbRestApiCallNodeTest extends AbstractRuleNodeUpgradeTest {
public void deleteRequestWithoutBody ( ) throws IOException , InterruptedException {
public void deleteRequestWithoutBody ( ) throws IOException , InterruptedException {
final CountDownLatch latch = new CountDownLatch ( 1 ) ;
final CountDownLatch latch = new CountDownLatch ( 1 ) ;
final String path = "/path/to/delete" ;
final String path = "/path/to/delete" ;
setupServer ( "*" , new HttpRequestHandler ( ) {
setupServer ( "*" , ( request , response , _ ) - > {
try {
@Override
assertEquals ( path , request . getRequestLine ( ) . getUri ( ) , "Request path matches" ) ;
public void handle ( HttpRequest request , HttpResponse response , HttpContext context )
assertTrue ( request . containsHeader ( "Foo" ) , "Custom header included" ) ;
throws HttpException , IOException {
assertEquals ( "Bar" , request . getFirstHeader ( "Foo" ) . getValue ( ) , "Custom header value" ) ;
try {
response . setStatusCode ( 200 ) ;
assertEquals ( request . getRequestLine ( ) . getUri ( ) , path , "Request path matches" ) ;
latch . countDown ( ) ;
assertTrue ( request . containsHeader ( "Foo" ) , "Custom header included" ) ;
} catch ( Exception e ) {
assertEquals ( "Bar" , request . getFirstHeader ( "Foo" ) . getValue ( ) , "Custom header value" ) ;
System . out . println ( "Exception handling request: " + e ) ;
response . setStatusCode ( 200 ) ;
e . printStackTrace ( ) ;
latch . countDown ( ) ;
latch . countDown ( ) ;
} catch ( Exception e ) {
System . out . println ( "Exception handling request: " + e . toString ( ) ) ;
e . printStackTrace ( ) ;
latch . countDown ( ) ;
}
}
}
} ) ;
} ) ;
@ -269,28 +261,23 @@ public class TbRestApiCallNodeTest extends AbstractRuleNodeUpgradeTest {
public void deleteRequestWithBody ( ) throws IOException , InterruptedException {
public void deleteRequestWithBody ( ) throws IOException , InterruptedException {
final CountDownLatch latch = new CountDownLatch ( 1 ) ;
final CountDownLatch latch = new CountDownLatch ( 1 ) ;
final String path = "/path/to/delete" ;
final String path = "/path/to/delete" ;
setupServer ( "*" , new HttpRequestHandler ( ) {
setupServer ( "*" , ( request , response , _ ) - > {
try {
@Override
assertEquals ( path , request . getRequestLine ( ) . getUri ( ) , "Request path matches" ) ;
public void handle ( HttpRequest request , HttpResponse response , HttpContext context )
assertTrue ( request . containsHeader ( "Content-Type" ) , "Content-Type included" ) ;
throws HttpException , IOException {
assertEquals ( "application/json" ,
try {
request . getFirstHeader ( "Content-Type" ) . getValue ( ) , "Content-Type value" ) ;
assertEquals ( path , request . getRequestLine ( ) . getUri ( ) , "Request path matches" ) ;
assertTrue ( request . containsHeader ( "Content-Length" ) , "Content-Length included" ) ;
assertTrue ( request . containsHeader ( "Content-Type" ) , "Content-Type included" ) ;
assertEquals ( "2" ,
assertEquals ( "application/json" ,
request . getFirstHeader ( "Content-Length" ) . getValue ( ) , "Content-Length value" ) ;
request . getFirstHeader ( "Content-Type" ) . getValue ( ) , "Content-Type value" ) ;
assertTrue ( request . containsHeader ( "Foo" ) , "Custom header included" ) ;
assertTrue ( request . containsHeader ( "Content-Length" ) , "Content-Length included" ) ;
assertEquals ( "Bar" , request . getFirstHeader ( "Foo" ) . getValue ( ) , "Custom header value" ) ;
assertEquals ( "2" ,
response . setStatusCode ( 200 ) ;
request . getFirstHeader ( "Content-Length" ) . getValue ( ) , "Content-Length value" ) ;
latch . countDown ( ) ;
assertTrue ( request . containsHeader ( "Foo" ) , "Custom header included" ) ;
} catch ( Exception e ) {
assertEquals ( "Bar" , request . getFirstHeader ( "Foo" ) . getValue ( ) , "Custom header value" ) ;
System . out . println ( "Exception handling request: " + e ) ;
response . setStatusCode ( 200 ) ;
e . printStackTrace ( ) ;
latch . countDown ( ) ;
latch . countDown ( ) ;
} catch ( Exception e ) {
System . out . println ( "Exception handling request: " + e . toString ( ) ) ;
e . printStackTrace ( ) ;
latch . countDown ( ) ;
}
}
}
} ) ;
} ) ;
@ -323,6 +310,143 @@ public class TbRestApiCallNodeTest extends AbstractRuleNodeUpgradeTest {
assertEquals ( TbMsg . EMPTY_JSON_OBJECT , dataCaptor . getValue ( ) ) ;
assertEquals ( TbMsg . EMPTY_JSON_OBJECT , dataCaptor . getValue ( ) ) ;
}
}
@Test
public void postRequestWithBodyTemplate ( ) throws IOException , InterruptedException {
final CountDownLatch latch = new CountDownLatch ( 1 ) ;
final String path = "/api/token" ;
final AtomicReference < String > capturedBody = new AtomicReference < > ( ) ;
setupServerWithBodyCapture ( capturedBody , latch ) ;
TbRestApiCallNodeConfiguration config = new TbRestApiCallNodeConfiguration ( ) . defaultConfiguration ( ) ;
config . setRequestMethod ( "POST" ) ;
config . setRequestBodyTemplate ( "{\"grant_type\":\"client_credentials\",\"client_id\":\"${clientId}\",\"value\":\"$[token]\"}" ) ;
config . setRestEndpointUrlPattern ( String . format ( "http://localhost:%d%s" , server . getLocalPort ( ) , path ) ) ;
initWithConfig ( config ) ;
metaData . putValue ( "clientId" , "my-client-123" ) ;
TbMsg msg = TbMsg . newMsg ( )
. type ( TbMsgType . POST_TELEMETRY_REQUEST )
. originator ( originator )
. copyMetaData ( metaData )
. dataType ( TbMsgDataType . JSON )
. data ( "{\"token\":\"abc-xyz\"}" )
. ruleChainId ( ruleChainId )
. ruleNodeId ( ruleNodeId )
. build ( ) ;
restNode . onMsg ( ctx , msg ) ;
assertTrue ( latch . await ( 10 , TimeUnit . SECONDS ) , "Server handled request" ) ;
assertEquals ( "{\"grant_type\":\"client_credentials\",\"client_id\":\"my-client-123\",\"value\":\"abc-xyz\"}" , capturedBody . get ( ) ) ;
}
@Test
public void postRequestWithBodyTemplateAndParseToPlainText ( ) throws IOException , InterruptedException {
final CountDownLatch latch = new CountDownLatch ( 1 ) ;
final String path = "/api/text" ;
final AtomicReference < String > capturedBody = new AtomicReference < > ( ) ;
setupServerWithBodyCapture ( capturedBody , latch ) ;
TbRestApiCallNodeConfiguration config = new TbRestApiCallNodeConfiguration ( ) . defaultConfiguration ( ) ;
config . setRequestMethod ( "POST" ) ;
config . setParseToPlainText ( true ) ;
config . setRequestBodyTemplate ( "Hello ${name}, your token is $[token]!" ) ;
config . setRestEndpointUrlPattern ( String . format ( "http://localhost:%d%s" , server . getLocalPort ( ) , path ) ) ;
initWithConfig ( config ) ;
metaData . putValue ( "name" , "World" ) ;
TbMsg msg = TbMsg . newMsg ( )
. type ( TbMsgType . POST_TELEMETRY_REQUEST )
. originator ( originator )
. copyMetaData ( metaData )
. dataType ( TbMsgDataType . JSON )
. data ( "{\"token\":\"abc-xyz\"}" )
. ruleChainId ( ruleChainId )
. ruleNodeId ( ruleNodeId )
. build ( ) ;
restNode . onMsg ( ctx , msg ) ;
assertTrue ( latch . await ( 10 , TimeUnit . SECONDS ) , "Server handled request" ) ;
assertEquals ( "Hello World, your token is abc-xyz!" , capturedBody . get ( ) ) ;
}
@Test
public void postRequestWithBodyTemplateEscapesJsonSpecialChars ( ) throws IOException , InterruptedException {
final CountDownLatch latch = new CountDownLatch ( 1 ) ;
final String path = "/api/token" ;
final AtomicReference < String > capturedBody = new AtomicReference < > ( ) ;
setupServerWithBodyCapture ( capturedBody , latch ) ;
TbRestApiCallNodeConfiguration config = new TbRestApiCallNodeConfiguration ( ) . defaultConfiguration ( ) ;
config . setRequestMethod ( "POST" ) ;
config . setRequestBodyTemplate ( "{\"name\":\"${userName}\",\"desc\":\"$[description]\"}" ) ;
config . setRestEndpointUrlPattern ( String . format ( "http://localhost:%d%s" , server . getLocalPort ( ) , path ) ) ;
initWithConfig ( config ) ;
metaData . putValue ( "userName" , "John \"Doe\"" ) ;
TbMsg msg = TbMsg . newMsg ( )
. type ( TbMsgType . POST_TELEMETRY_REQUEST )
. originator ( originator )
. copyMetaData ( metaData )
. dataType ( TbMsgDataType . JSON )
. data ( "{\"description\":\"line1\\nline2\"}" )
. ruleChainId ( ruleChainId )
. ruleNodeId ( ruleNodeId )
. build ( ) ;
restNode . onMsg ( ctx , msg ) ;
assertTrue ( latch . await ( 10 , TimeUnit . SECONDS ) , "Server handled request" ) ;
assertEquals ( "{\"name\":\"John \\\"Doe\\\"\",\"desc\":\"line1\\nline2\"}" , capturedBody . get ( ) ) ;
}
@Test
public void postRequestWithEmptyBodyTemplateUsesMessageData ( ) throws IOException , InterruptedException {
final CountDownLatch latch = new CountDownLatch ( 1 ) ;
final String path = "/api/data" ;
final AtomicReference < String > capturedBody = new AtomicReference < > ( ) ;
setupServerWithBodyCapture ( capturedBody , latch ) ;
TbRestApiCallNodeConfiguration config = new TbRestApiCallNodeConfiguration ( ) . defaultConfiguration ( ) ;
config . setRequestMethod ( "POST" ) ;
// requestBodyTemplate is null by default — should use msg.getData()
config . setRestEndpointUrlPattern ( String . format ( "http://localhost:%d%s" , server . getLocalPort ( ) , path ) ) ;
initWithConfig ( config ) ;
TbMsg msg = TbMsg . newMsg ( )
. type ( TbMsgType . POST_TELEMETRY_REQUEST )
. originator ( originator )
. copyMetaData ( metaData )
. dataType ( TbMsgDataType . JSON )
. data ( "{\"temperature\":25}" )
. ruleChainId ( ruleChainId )
. ruleNodeId ( ruleNodeId )
. build ( ) ;
restNode . onMsg ( ctx , msg ) ;
assertTrue ( latch . await ( 10 , TimeUnit . SECONDS ) , "Server handled request" ) ;
ArgumentCaptor < TbMsg > msgCaptor = ArgumentCaptor . forClass ( TbMsg . class ) ;
ArgumentCaptor < TbMsgMetaData > metadataCaptor = ArgumentCaptor . forClass ( TbMsgMetaData . class ) ;
ArgumentCaptor < String > dataCaptor = ArgumentCaptor . forClass ( String . class ) ;
verify ( ctx , timeout ( 10_000 ) ) . transformMsg ( msgCaptor . capture ( ) , metadataCaptor . capture ( ) , dataCaptor . capture ( ) ) ;
assertEquals ( "{\"temperature\":25}" , capturedBody . get ( ) ) ;
}
private void setupServerWithBodyCapture ( AtomicReference < String > capturedBody , CountDownLatch latch ) throws IOException {
setupServer ( "*" , ( request , response , _ ) - > {
try {
if ( request instanceof org . apache . http . HttpEntityEnclosingRequest entityRequest ) {
InputStream is = entityRequest . getEntity ( ) . getContent ( ) ;
capturedBody . set ( new String ( is . readAllBytes ( ) , StandardCharsets . UTF_8 ) ) ;
}
response . setStatusCode ( 200 ) ;
latch . countDown ( ) ;
} catch ( Exception e ) {
e . printStackTrace ( ) ;
latch . countDown ( ) ;
}
} ) ;
}
private static Stream < Arguments > givenFromVersionAndConfig_whenUpgrade_thenVerifyHasChangesAndConfig ( ) {
private static Stream < Arguments > givenFromVersionAndConfig_whenUpgrade_thenVerifyHasChangesAndConfig ( ) {
return Stream . of (
return Stream . of (
Arguments . of ( 0 ,
Arguments . of ( 0 ,
@ -339,7 +463,7 @@ public class TbRestApiCallNodeTest extends AbstractRuleNodeUpgradeTest {
"\"proxyPort\": 0,\"proxyUser\": null,\"proxyPassword\": null,\"readTimeoutMs\": 0," +
"\"proxyPort\": 0,\"proxyUser\": null,\"proxyPassword\": null,\"readTimeoutMs\": 0," +
"\"maxParallelRequestsCount\": 0,\"headers\": {\"Content-Type\": \"application/json\"}," +
"\"maxParallelRequestsCount\": 0,\"headers\": {\"Content-Type\": \"application/json\"}," +
"\"credentials\": {\"type\": \"anonymous\"}," +
"\"credentials\": {\"type\": \"anonymous\"}," +
"\"maxInMemoryBufferSizeInKb\": 256}" ) ,
"\"maxInMemoryBufferSizeInKb\": 256,\"requestBodyTemplate\": null }" ) ,
Arguments . of ( 1 ,
Arguments . of ( 1 ,
"{\"restEndpointUrlPattern\":\"http://localhost/api\",\"requestMethod\": \"POST\"," +
"{\"restEndpointUrlPattern\":\"http://localhost/api\",\"requestMethod\": \"POST\"," +
"\"useSimpleClientHttpFactory\": false,\"parseToPlainText\": false,\"ignoreRequestBody\": false," +
"\"useSimpleClientHttpFactory\": false,\"parseToPlainText\": false,\"ignoreRequestBody\": false," +
@ -355,7 +479,7 @@ public class TbRestApiCallNodeTest extends AbstractRuleNodeUpgradeTest {
"\"proxyPort\": 0,\"proxyUser\": null,\"proxyPassword\": null,\"readTimeoutMs\": 0," +
"\"proxyPort\": 0,\"proxyUser\": null,\"proxyPassword\": null,\"readTimeoutMs\": 0," +
"\"maxParallelRequestsCount\": 0,\"headers\": {\"Content-Type\": \"application/json\"}," +
"\"maxParallelRequestsCount\": 0,\"headers\": {\"Content-Type\": \"application/json\"}," +
"\"credentials\": {\"type\": \"anonymous\"}," +
"\"credentials\": {\"type\": \"anonymous\"}," +
"\"maxInMemoryBufferSizeInKb\": 256}" ) ,
"\"maxInMemoryBufferSizeInKb\": 256,\"requestBodyTemplate\": null }" ) ,
Arguments . of ( 2 ,
Arguments . of ( 2 ,
"{\"restEndpointUrlPattern\":\"http://localhost/api\",\"requestMethod\": \"POST\"," +
"{\"restEndpointUrlPattern\":\"http://localhost/api\",\"requestMethod\": \"POST\"," +
"\"useSimpleClientHttpFactory\": false,\"parseToPlainText\": false,\"ignoreRequestBody\": false," +
"\"useSimpleClientHttpFactory\": false,\"parseToPlainText\": false,\"ignoreRequestBody\": false," +
@ -370,7 +494,7 @@ public class TbRestApiCallNodeTest extends AbstractRuleNodeUpgradeTest {
"\"proxyPort\": 0,\"proxyUser\": null,\"proxyPassword\": null,\"readTimeoutMs\": 0," +
"\"proxyPort\": 0,\"proxyUser\": null,\"proxyPassword\": null,\"readTimeoutMs\": 0," +
"\"maxParallelRequestsCount\": 0,\"headers\": {\"Content-Type\": \"application/json\"}," +
"\"maxParallelRequestsCount\": 0,\"headers\": {\"Content-Type\": \"application/json\"}," +
"\"credentials\": {\"type\": \"anonymous\"}," +
"\"credentials\": {\"type\": \"anonymous\"}," +
"\"maxInMemoryBufferSizeInKb\": 256}" ) ,
"\"maxInMemoryBufferSizeInKb\": 256,\"requestBodyTemplate\": null }" ) ,
Arguments . of ( 3 , "" "
Arguments . of ( 3 , "" "
{
{
"restEndpointUrlPattern" : "http://localhost/api" ,
"restEndpointUrlPattern" : "http://localhost/api" ,
@ -419,6 +543,7 @@ public class TbRestApiCallNodeTest extends AbstractRuleNodeUpgradeTest {
"credentials" : {
"credentials" : {
"type" : "anonymous"
"type" : "anonymous"
} ,
} ,
"requestBodyTemplate" : null ,
"maxInMemoryBufferSizeInKb" : 256
"maxInMemoryBufferSizeInKb" : 256
} "" " )
} "" " )
) ;
) ;