Browse Source

REST API call node: better URL encoding

# Conflicts:
#	common/data/src/main/java/org/thingsboard/server/common/data/util/CollectionsUtil.java
pull/14540/head
Dmytro Skarzhynets 10 months ago
parent
commit
3a098ef815
No known key found for this signature in database GPG Key ID: 2B51652F224037DF
  1. 10
      common/data/src/main/java/org/thingsboard/server/common/data/util/CollectionsUtil.java
  2. 4
      rule-engine/rule-engine-components/pom.xml
  3. 58
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbHttpClient.java
  4. 13
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbRestApiCallNode.java
  5. 26
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbRestApiCallNodeConfiguration.java
  6. 498
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/rest/TbHttpClientTest.java
  7. 114
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/rest/TbRestApiCallNodeTest.java

10
common/data/src/main/java/org/thingsboard/server/common/data/util/CollectionsUtil.java

@ -18,6 +18,7 @@ package org.thingsboard.server.common.data.util;
import java.util.Collection; import java.util.Collection;
import java.util.HashMap; import java.util.HashMap;
import java.util.HashSet; import java.util.HashSet;
import java.util.LinkedHashMap;
import java.util.List; import java.util.List;
import java.util.Map; import java.util.Map;
import java.util.Set; import java.util.Set;
@ -72,6 +73,15 @@ public class CollectionsUtil {
return map; return map;
} }
@SafeVarargs
public static <K, V> LinkedHashMap<K, V> orderedMapOf(Map.Entry<K, V>... entries) {
LinkedHashMap<K, V> map = new LinkedHashMap<>();
for (Map.Entry<K, V> entry : entries) {
map.put(entry.getKey(), entry.getValue());
}
return map;
}
public static <V> boolean emptyOrContains(Collection<V> collection, V element) { public static <V> boolean emptyOrContains(Collection<V> collection, V element) {
return isEmpty(collection) || collection.contains(element); return isEmpty(collection) || collection.contains(element);
} }

4
rule-engine/rule-engine-components/pom.xml

@ -161,6 +161,10 @@
<groupId>org.thingsboard.langchain4j</groupId> <groupId>org.thingsboard.langchain4j</groupId>
<artifactId>langchain4j</artifactId> <artifactId>langchain4j</artifactId>
</dependency> </dependency>
<dependency>
<groupId>org.apache.httpcomponents.core5</groupId>
<artifactId>httpcore5</artifactId>
</dependency>
</dependencies> </dependencies>
<build> <build>

58
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbHttpClient.java

@ -15,12 +15,14 @@
*/ */
package org.thingsboard.rule.engine.rest; package org.thingsboard.rule.engine.rest;
import com.google.common.collect.Maps;
import io.netty.channel.EventLoopGroup; import io.netty.channel.EventLoopGroup;
import io.netty.channel.nio.NioEventLoopGroup; import io.netty.channel.nio.NioEventLoopGroup;
import io.netty.handler.ssl.SslContext; import io.netty.handler.ssl.SslContext;
import io.netty.handler.timeout.ReadTimeoutHandler; import io.netty.handler.timeout.ReadTimeoutHandler;
import lombok.Data; import lombok.Data;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.apache.hc.core5.net.URIBuilder;
import org.springframework.http.HttpHeaders; import org.springframework.http.HttpHeaders;
import org.springframework.http.HttpMethod; import org.springframework.http.HttpMethod;
import org.springframework.http.HttpStatus; import org.springframework.http.HttpStatus;
@ -47,6 +49,7 @@ import reactor.netty.transport.ProxyProvider;
import javax.net.ssl.SSLException; import javax.net.ssl.SSLException;
import java.net.URI; import java.net.URI;
import java.net.URISyntaxException;
import java.nio.charset.StandardCharsets; import java.nio.charset.StandardCharsets;
import java.util.Base64; import java.util.Base64;
import java.util.List; import java.util.List;
@ -213,8 +216,21 @@ public class TbHttpClient {
} }
String endpointUrl = TbNodeUtils.processPattern(config.getRestEndpointUrlPattern(), msg); String endpointUrl = TbNodeUtils.processPattern(config.getRestEndpointUrlPattern(), msg);
Map<String, String> processedQueryParams;
if (config.isUseNewEncoding()) {
processedQueryParams = Maps.newHashMapWithExpectedSize(config.getQueryParams().size());
config.getQueryParams().forEach((name, value) -> {
var processedParamName = TbNodeUtils.processPattern(name, msg);
var processedParamValue = TbNodeUtils.processPattern(value, msg);
processedQueryParams.put(processedParamName, processedParamValue);
});
} else {
processedQueryParams = null;
}
HttpMethod method = HttpMethod.valueOf(config.getRequestMethod()); HttpMethod method = HttpMethod.valueOf(config.getRequestMethod());
URI uri = buildEncodedUri(endpointUrl); URI uri = buildEncodedUri(endpointUrl, processedQueryParams);
RequestBodySpec request = webClient RequestBodySpec request = webClient
.method(method) .method(method)
@ -262,7 +278,7 @@ public class TbHttpClient {
return origin; return origin;
} }
public URI buildEncodedUri(String endpointUrl) { public URI buildEncodedUri(String endpointUrl, Map<String, String> queryParams) {
if (endpointUrl == null) { if (endpointUrl == null) {
throw new RuntimeException("Url string cannot be null!"); throw new RuntimeException("Url string cannot be null!");
} }
@ -270,7 +286,13 @@ public class TbHttpClient {
throw new RuntimeException("Url string cannot be empty!"); throw new RuntimeException("Url string cannot be empty!");
} }
URI uri = UriComponentsBuilder.fromUriString(endpointUrl).build().encode().toUri(); URI uri;
if (queryParams != null) {
uri = buildEncodedUriNew(endpointUrl, queryParams);
} else {
uri = buildEncodedUriLegacy(endpointUrl);
}
if (uri.getScheme() == null || uri.getScheme().isEmpty()) { if (uri.getScheme() == null || uri.getScheme().isEmpty()) {
throw new RuntimeException("Transport scheme(protocol) must be provided!"); throw new RuntimeException("Transport scheme(protocol) must be provided!");
} }
@ -284,6 +306,36 @@ public class TbHttpClient {
return uri; return uri;
} }
private URI buildEncodedUriNew(String endpointUrl, Map<String, String> queryParams) {
try {
URIBuilder builder = new URIBuilder(endpointUrl);
queryParams.forEach(builder::addParameter);
return builder.build();
} catch (URISyntaxException e) {
throw new IllegalArgumentException("""
Invalid REST endpoint URL: '%s'. The URL must be valid and properly encoded.
If the path contains special characters (e.g., spaces, braces), they must be percent-encoded (e.g., '/my file/' should be '/my%%20file/', '/path/{id}/' should be '/path/%%7Bid%%7D/').
Query parameters should be configured separately and will be encoded automatically.""".formatted(endpointUrl), e);
}
}
/**
* @deprecated Query params embedded in the URL string are not encoded correctly.
* <p>
* In query strings, {@code +} is interpreted as a space (URL form encoding),
* so {@code email=user+tag@test.com} becomes {@code email=user tag@test.com} on the server.
* <p>
* Pre-encoding the URL (e.g., {@code email=user%2Btag@test.com}) doesn't help
* because {@code .encode()} will double-encode it to {@code email=user%252Btag@test.com}.
* <p>
* Use {@link #buildEncodedUriNew} with a separate query params map,
* where values are encoded exactly once: {@code +} → {@code %2B}, {@code @} → {@code %40}.
*/
@Deprecated(since = "4.2.1.2", forRemoval = true) // fixme: remove in major 5.0
private URI buildEncodedUriLegacy(String endpointUrl) {
return UriComponentsBuilder.fromUriString(endpointUrl).build().encode().toUri();
}
private Object getData(TbMsg tbMsg, boolean parseToPlainText) { private Object getData(TbMsg tbMsg, boolean parseToPlainText) {
String data = tbMsg.getData(); String data = tbMsg.getData();
return parseToPlainText ? JacksonUtil.toPlainText(data) : JacksonUtil.toJsonNode(data); return parseToPlainText ? JacksonUtil.toPlainText(data) : JacksonUtil.toJsonNode(data);

13
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbRestApiCallNode.java

@ -26,9 +26,12 @@ import org.thingsboard.rule.engine.external.TbAbstractExternalNode;
import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.data.util.TbPair; import org.thingsboard.server.common.data.util.TbPair;
import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.dao.exception.DataValidationException;
import java.util.List; import java.util.List;
import static org.thingsboard.server.dao.service.ConstraintValidator.validateFields;
@RuleNode( @RuleNode(
type = ComponentType.EXTERNAL, type = ComponentType.EXTERNAL,
name = "rest api call", name = "rest api call",
@ -57,7 +60,15 @@ public class TbRestApiCallNode extends TbAbstractExternalNode {
@Override @Override
public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException { public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException {
super.init(ctx); super.init(ctx);
TbRestApiCallNodeConfiguration config = TbNodeUtils.convert(configuration, TbRestApiCallNodeConfiguration.class);
var config = TbNodeUtils.convert(configuration, TbRestApiCallNodeConfiguration.class);
String errorPrefix = "'" + ctx.getSelf().getName() + "' node configuration is invalid: ";
try {
validateFields(config, errorPrefix);
} catch (DataValidationException e) {
throw new TbNodeException(e, true);
}
httpClient = new TbHttpClient(config, ctx.getSharedEventLoop()); httpClient = new TbHttpClient(config, ctx.getSharedEventLoop());
} }

26
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbRestApiCallNodeConfiguration.java

@ -15,7 +15,10 @@
*/ */
package org.thingsboard.rule.engine.rest; package org.thingsboard.rule.engine.rest;
import com.fasterxml.jackson.annotation.JsonIgnore;
import com.fasterxml.jackson.annotation.JsonIgnoreProperties; import com.fasterxml.jackson.annotation.JsonIgnoreProperties;
import jakarta.validation.constraints.AssertTrue;
import jakarta.validation.constraints.NotNull;
import lombok.Data; import lombok.Data;
import org.springframework.http.HttpHeaders; import org.springframework.http.HttpHeaders;
import org.springframework.http.MediaType; import org.springframework.http.MediaType;
@ -32,6 +35,8 @@ public class TbRestApiCallNodeConfiguration implements NodeConfiguration<TbRestA
private String restEndpointUrlPattern; private String restEndpointUrlPattern;
private String requestMethod; private String requestMethod;
private boolean useNewEncoding;
private Map<@NotNull String, @NotNull(message = "query parameter values must be non-null") String> queryParams;
private Map<String, String> headers; private Map<String, String> headers;
private boolean useSimpleClientHttpFactory; private boolean useSimpleClientHttpFactory;
private int readTimeoutMs; private int readTimeoutMs;
@ -48,12 +53,32 @@ public class TbRestApiCallNodeConfiguration implements NodeConfiguration<TbRestA
private boolean ignoreRequestBody; private boolean ignoreRequestBody;
private int maxInMemoryBufferSizeInKb; private int maxInMemoryBufferSizeInKb;
@JsonIgnore
@AssertTrue(message = "query parameters must be non-null if new encoding is used")
public boolean isNewEncodingConfigValid() {
if (useNewEncoding) {
return queryParams != null;
}
return true;
}
@JsonIgnore
@AssertTrue(message = "query parameters must be null if old encoding is used")
public boolean isLegacyEncodingConfigValid() {
if (!useNewEncoding) {
return queryParams == null;
}
return true;
}
@Override @Override
public TbRestApiCallNodeConfiguration defaultConfiguration() { public TbRestApiCallNodeConfiguration defaultConfiguration() {
TbRestApiCallNodeConfiguration configuration = new TbRestApiCallNodeConfiguration(); TbRestApiCallNodeConfiguration configuration = new TbRestApiCallNodeConfiguration();
configuration.setRestEndpointUrlPattern("http://localhost/api"); configuration.setRestEndpointUrlPattern("http://localhost/api");
configuration.setRequestMethod("POST"); configuration.setRequestMethod("POST");
configuration.setHeaders(Collections.singletonMap(HttpHeaders.CONTENT_TYPE, MediaType.APPLICATION_JSON_VALUE)); configuration.setHeaders(Collections.singletonMap(HttpHeaders.CONTENT_TYPE, MediaType.APPLICATION_JSON_VALUE));
configuration.setUseNewEncoding(true);
configuration.setQueryParams(Collections.emptyMap());
configuration.setUseSimpleClientHttpFactory(false); configuration.setUseSimpleClientHttpFactory(false);
configuration.setReadTimeoutMs(0); configuration.setReadTimeoutMs(0);
configuration.setMaxParallelRequestsCount(0); configuration.setMaxParallelRequestsCount(0);
@ -72,4 +97,5 @@ public class TbRestApiCallNodeConfiguration implements NodeConfiguration<TbRestA
return this.credentials; return this.credentials;
} }
} }
} }

498
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/rest/TbHttpClientTest.java

@ -15,16 +15,17 @@
*/ */
package org.thingsboard.rule.engine.rest; package org.thingsboard.rule.engine.rest;
import io.netty.channel.EventLoopGroup; import io.netty.channel.EventLoopGroup;
import io.netty.channel.nio.NioEventLoopGroup; import io.netty.channel.nio.NioEventLoopGroup;
import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Named;
import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Test;
import org.mockito.ArgumentCaptor; import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.Arguments;
import org.junit.jupiter.params.provider.MethodSource;
import org.mockito.Mockito; import org.mockito.Mockito;
import org.mockserver.integration.ClientAndServer;
import org.springframework.util.LinkedMultiValueMap; import org.springframework.util.LinkedMultiValueMap;
import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.DeviceId;
@ -34,25 +35,33 @@ import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgMetaData; import org.thingsboard.server.common.msg.TbMsgMetaData;
import java.net.URI; import java.net.URI;
import java.util.Collections;
import java.util.List; import java.util.List;
import java.util.Map; import java.util.Map;
import java.util.UUID;
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 static java.util.Map.entry;
import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.hamcrest.MatcherAssert.assertThat; import static org.hamcrest.MatcherAssert.assertThat;
import static org.hamcrest.Matchers.instanceOf; import static org.hamcrest.Matchers.instanceOf;
import static org.hamcrest.Matchers.is; import static org.hamcrest.Matchers.is;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.eq; import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.BDDMockito.willCallRealMethod; import static org.mockito.BDDMockito.willCallRealMethod;
import static org.mockito.Mockito.mock; import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.times; import static org.mockito.Mockito.never;
import static org.mockito.Mockito.verify; import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when; import static org.mockito.Mockito.when;
import static org.mockserver.integration.ClientAndServer.startClientAndServer; import static org.mockserver.integration.ClientAndServer.startClientAndServer;
import static org.mockserver.model.HttpRequest.request; import static org.mockserver.model.HttpRequest.request;
import static org.mockserver.model.HttpResponse.response; import static org.mockserver.model.HttpResponse.response;
import static org.thingsboard.server.common.data.util.CollectionsUtil.orderedMapOf;
public class TbHttpClientTest { public class TbHttpClientTest {
@ -86,123 +95,456 @@ public class TbHttpClientTest {
@Test @Test
public void testBuildSimpleUri() { public void testBuildSimpleUri() {
Mockito.when(client.buildEncodedUri(any())).thenCallRealMethod(); Mockito.when(client.buildEncodedUri(any(), any())).thenCallRealMethod();
String url = "http://localhost:8080/"; String url = "http://localhost:8080/";
URI uri = client.buildEncodedUri(url); URI uri = client.buildEncodedUri(url, null);
Assertions.assertEquals(url, uri.toString()); Assertions.assertEquals(url, uri.toString());
} }
@Test @Test
public void testBuildUriWithoutProtocol() { public void testBuildUriWithoutProtocol() {
Mockito.when(client.buildEncodedUri(any())).thenCallRealMethod(); Mockito.when(client.buildEncodedUri(any(), any())).thenCallRealMethod();
String url = "localhost:8080/"; String url = "localhost:8080/";
assertThatThrownBy(() -> client.buildEncodedUri(url)); assertThatThrownBy(() -> client.buildEncodedUri(url, null));
} }
@Test @Test
public void testBuildInvalidUri() { public void testBuildInvalidUri() {
Mockito.when(client.buildEncodedUri(any())).thenCallRealMethod(); Mockito.when(client.buildEncodedUri(any(), any())).thenCallRealMethod();
String url = "aaa"; String url = "aaa";
assertThatThrownBy(() -> client.buildEncodedUri(url)); assertThatThrownBy(() -> client.buildEncodedUri(url, null));
} }
@Test @Test
public void testBuildUriWithSpecialSymbols() { public void testBuildUriWithSpecialSymbols() {
Mockito.when(client.buildEncodedUri(any())).thenCallRealMethod(); Mockito.when(client.buildEncodedUri(any(), any())).thenCallRealMethod();
String url = "http://192.168.1.1/data?d={\"a\": 12}"; String url = "http://192.168.1.1/data?d={\"a\": 12}";
String expected = "http://192.168.1.1/data?d=%7B%22a%22:%2012%7D"; String expected = "http://192.168.1.1/data?d=%7B%22a%22:%2012%7D";
URI uri = client.buildEncodedUri(url); URI uri = client.buildEncodedUri(url, null);
Assertions.assertEquals(expected, uri.toString()); Assertions.assertEquals(expected, uri.toString());
} }
@Test
public void testProcessMessageWithQueryParamsPatternProcessing() throws Exception {
// GIVEN
String path = "/api/notify";
try (var server = startClientAndServer("localhost", 0)) {
server.when(
request()
.withMethod("GET")
.withPath(path)
.withQueryStringParameter("email", "user+tag@test.com")
.withQueryStringParameter("device", "sensor-1")
.withQueryStringParameter("custom-header", "header-value")
.withQueryStringParameter("temp", "25.5")
.withQueryStringParameter("location", "room-1/zone-A")
).respond(
response().withStatusCode(200)
);
var config = new TbRestApiCallNodeConfiguration().defaultConfiguration();
config.setRestEndpointUrlPattern("http://localhost:" + server.getPort() + path);
config.setRequestMethod("GET");
config.setUseSimpleClientHttpFactory(true);
config.setUseNewEncoding(true);
config.setQueryParams(Map.of(
"email", "${userEmail}", // ${} from metadata
"device", "${deviceName}", // ${} from metadata
"${dynamicParam}", "${dynamicValue}", // ${} in both key and value
"temp", "$[temperature]", // $[] from data
"location", "$[sensor.location]" // $[] from nested data
));
var metaData = new TbMsgMetaData();
metaData.putValue("userEmail", "user+tag@test.com");
metaData.putValue("deviceName", "sensor-1");
metaData.putValue("dynamicParam", "custom-header");
metaData.putValue("dynamicValue", "header-value");
var msg = TbMsg.newMsg()
.type(TbMsgType.POST_TELEMETRY_REQUEST)
.originator(new DeviceId(EntityId.NULL_UUID))
.metaData(metaData)
.data("""
{"temperature": 25.5, "sensor": {"location": "room-1/zone-A"}}
""")
.build();
// WHEN-THEN
processMessageAndWait(config, msg);
}
}
@Test @Test
public void testProcessMessageWithJsonInUrlVariable() throws Exception { public void testProcessMessageWithJsonInUrlVariable() throws Exception {
String host = "localhost"; // GIVEN
String path = "/api"; String path = "/api";
String paramKey = "data"; String paramValue = "[{\"test\":\"test\"}]";
String paramVal = "[{\"test\":\"test\"}]";
String successResponseBody = "SUCCESS"; try (var server = startClientAndServer("localhost", 0)) {
server.when(
var server = setUpDummyServer(host, path, paramKey, paramVal, successResponseBody); request()
.withMethod("GET")
String endpointUrl = String.format( .withPath(path)
"http://%s:%d%s?%s=%s", .withQueryStringParameter("data", paramValue)
host, server.getPort(), path, paramKey, paramVal ).respond(
); response().withStatusCode(200)
String method = "GET"; );
var config = new TbRestApiCallNodeConfiguration().defaultConfiguration();
var config = new TbRestApiCallNodeConfiguration() config.setRestEndpointUrlPattern("http://localhost:" + server.getPort() + path + "?data=" + paramValue);
.defaultConfiguration(); config.setRequestMethod("GET");
config.setRequestMethod(method); config.setUseNewEncoding(false);
config.setRestEndpointUrlPattern(endpointUrl); config.setQueryParams(null);
config.setUseSimpleClientHttpFactory(true);
var msg = TbMsg.newMsg()
.type(TbMsgType.POST_TELEMETRY_REQUEST)
.originator(new DeviceId(UUID.randomUUID()))
.metaData(new TbMsgMetaData())
.data(TbMsg.EMPTY_JSON_OBJECT)
.build();
// WHEN-THEN
processMessageAndWait(config, msg);
}
}
private void processMessageAndWait(TbRestApiCallNodeConfiguration config, TbMsg msg) throws Exception {
var httpClient = new TbHttpClient(config, eventLoop); var httpClient = new TbHttpClient(config, eventLoop);
var msg = TbMsg.newMsg()
.type(TbMsgType.POST_TELEMETRY_REQUEST)
.originator(new DeviceId(EntityId.NULL_UUID))
.copyMetaData(TbMsgMetaData.EMPTY)
.data(TbMsg.EMPTY_JSON_OBJECT)
.build();
var successMsg = TbMsg.newMsg()
.type(TbMsgType.POST_TELEMETRY_REQUEST)
.originator(msg.getOriginator())
.copyMetaData(msg.getMetaData())
.data(msg.getData())
.build();
var ctx = mock(TbContext.class); var ctx = mock(TbContext.class);
when(ctx.transformMsg( when(ctx.transformMsg(eq(msg), any(), any())).thenReturn(msg);
eq(msg),
eq(msg.getMetaData()),
eq(msg.getData())
)).thenReturn(successMsg);
var capturedData = ArgumentCaptor.forClass(String.class);
when(ctx.transformMsg( var latch = new CountDownLatch(1);
eq(msg), var error = new AtomicReference<Throwable>();
any(),
capturedData.capture()
)).thenReturn(successMsg);
CountDownLatch latch = new CountDownLatch(1);
httpClient.processMessage(ctx, msg, httpClient.processMessage(ctx, msg,
m -> { m -> {
ctx.tellSuccess(msg); ctx.tellSuccess(m);
latch.countDown(); latch.countDown();
}, },
(m, t) -> { (m, t) -> {
ctx.tellFailure(m, t); ctx.tellFailure(m, t);
error.set(t);
latch.countDown(); latch.countDown();
}); });
latch.await(5, TimeUnit.SECONDS); assertTrue(latch.await(5, TimeUnit.SECONDS), "Request should complete within timeout");
assertNull(error.get(), "Request should succeed, but got: " + error.get());
verify(ctx, times(1)).tellSuccess(any()); verify(ctx).tellSuccess(any());
verify(ctx, times(0)).tellFailure(any(), any()); verify(ctx, never()).tellFailure(any(), any());
Assertions.assertEquals(successResponseBody, capturedData.getValue());
} }
private ClientAndServer setUpDummyServer(String host, String path, String paramKey, String paramVal, String successResponseBody) {
var server = startClientAndServer(host, 1080); @ParameterizedTest(name = "{0}")
createGetMethodExpectations(server, path, paramKey, paramVal, successResponseBody); @MethodSource
return server; public void testQueryParamsEncoding(String endpointUrl, Map<String, String> queryParams, String expectedEncodedUrl) {
// GIVEN
Mockito.when(client.buildEncodedUri(any(), any())).thenCallRealMethod();
// WHEN
URI uri = client.buildEncodedUri(endpointUrl, queryParams);
// THEN
Assertions.assertEquals(expectedEncodedUrl, uri.toASCIIString());
} }
private void createGetMethodExpectations(ClientAndServer server, String path, String paramKey, String paramVal, String successResponseBody) { private static Stream<Arguments> testQueryParamsEncoding() {
server.when( return Stream.of(
request() Arguments.of(
.withMethod("GET") Named.named("ISO 8601 date-time in value", "http://somecompany/api/data/fetch"),
.withPath(path) Map.of("ts", "2016-08-01T09:06:06.0+02:00"),
.withQueryStringParameter(paramKey, paramVal) "http://somecompany/api/data/fetch?ts=2016-08-01T09%3A06%3A06.0%2B02%3A00"
).respond( ),
response() Arguments.of(
.withStatusCode(200) Named.named("email with plus sign in value", "http://localhost:8080/api/user/sendActivationMail"),
.withBody(successResponseBody) Map.of("email", "someperson+test1289@thingsboard.io"),
"http://localhost:8080/api/user/sendActivationMail?email=someperson%2Btest1289%40thingsboard.io"
),
Arguments.of(
Named.named("plus mixed with spaces in value", "http://url/api"),
Map.of("q", "a + b"),
"http://url/api?q=a%20%2B%20b"
),
Arguments.of(
Named.named("colon in value", "http://url/api"),
Map.of("time", "12:00"),
"http://url/api?time=12%3A00"
),
Arguments.of(
Named.named("slash in value", "http://url"),
Map.of("ref", "/home/user"),
"http://url?ref=%2Fhome%2Fuser"
),
Arguments.of(
Named.named("comma and semicolon in value", "http://url"),
Map.of("l", "a,b;c"),
"http://url?l=a%2Cb%3Bc"
),
Arguments.of(
Named.named("ampersand and equals in value", "http://url"),
Map.of("q", "key1=value1&key2=value2"),
"http://url?q=key1%3Dvalue1%26key2%3Dvalue2"
),
Arguments.of(
Named.named("JSON in value", "http://url"),
Map.of("json", """
{
"string": "hello",
"integer": 42,
"float": 3.14,
"boolTrue": true,
"boolFalse": false,
"null": null,
"array": [
1,
"two",
true,
null
],
"object": {
"nested": "value"
}
}"""),
"http://url?json=%7B%0A%20%20%20%20%22string%22%3A%20%22hello%22%2C%0A%20%20%20%20%22integer%22%3A%2042%2C%0A%20%20%20%20%22float%22%3A%203.14%2C%0A%20%20%20%20%22boolTrue%22%3A%20true%2C%0A%20%20%20%20%22boolFalse%22%3A%20false%2C%0A%20%20%20%20%22null%22%3A%20null%2C%0A%20%20%20%20%22array%22%3A%20%5B%0A%20%20%20%20%20%20%20%201%2C%0A%20%20%20%20%20%20%20%20%22two%22%2C%0A%20%20%20%20%20%20%20%20true%2C%0A%20%20%20%20%20%20%20%20null%0A%20%20%20%20%5D%2C%0A%20%20%20%20%22object%22%3A%20%7B%0A%20%20%20%20%20%20%20%20%22nested%22%3A%20%22value%22%0A%20%20%20%20%7D%0A%7D"
),
Arguments.of(
Named.named("UTF-8 in query", "http://url/cafes"),
orderedMapOf(
entry("nom", "Le Goût Moderne"),
entry("назва", "У Миколи \uD83D\uDE0B")
),
"http://url/cafes?nom=Le%20Go%C3%BBt%20Moderne&%D0%BD%D0%B0%D0%B7%D0%B2%D0%B0=%D0%A3%20%D0%9C%D0%B8%D0%BA%D0%BE%D0%BB%D0%B8%20%F0%9F%98%8B"
),
Arguments.of(
Named.named("empty value", "http://url/empty"),
Map.of("name", ""),
"http://url/empty?name="
),
Arguments.of(
Named.named("blank value (spaces)", "http://url/empty"),
Map.of("name", " "),
"http://url/empty?name=%20%20"
),
Arguments.of(
Named.named("empty key", "http://url/empty"),
Map.of("", "value"),
"http://url/empty?=value"
),
Arguments.of(
Named.named("blank key (spaces)", "http://url/empty"),
Map.of(" ", "value"),
"http://url/empty?%20%20=value"
),
Arguments.of(
Named.named("blank key with value that needs to be encoded", "http://url/empty"),
Map.of(" ", "value1+value2"),
"http://url/empty?%20%20=value1%2Bvalue2"
),
Arguments.of(
Named.named("blank key and value", "http://url/empty"),
Map.of(" ", " "),
"http://url/empty?%20%20=%20%20"
),
Arguments.of(
Named.named("fragment with query params", "http://url#frag"),
Map.of("docs", "مستندات عقدة القاعدة\n"),
"http://url?docs=%D9%85%D8%B3%D8%AA%D9%86%D8%AF%D8%A7%D8%AA%20%D8%B9%D9%82%D8%AF%D8%A9%20%D8%A7%D9%84%D9%82%D8%A7%D8%B9%D8%AF%D8%A9%0A#frag"
),
Arguments.of(
Named.named("fragment only", "http://url#frag"),
Collections.emptyMap(),
"http://url#frag"
),
Arguments.of(
Named.named("multiple query params (ordered)", "http://url/api"),
orderedMapOf(
entry("param1", "value1"),
entry("param2", "value2"),
entry("param3", "value3")
),
"http://url/api?param1=value1&param2=value2&param3=value3"
),
Arguments.of(
Named.named("hash (#) in value", "http://url/api"),
Map.of("color", "#ff0000"),
"http://url/api?color=%23ff0000"
),
Arguments.of(
Named.named("question mark (?) in value", "http://url/api"),
Map.of("query", "what?"),
"http://url/api?query=what%3F"
),
Arguments.of(
Named.named("percent sign (%) in value", "http://url/api"),
Map.of("discount", "50%"),
"http://url/api?discount=50%25"
),
Arguments.of(
Named.named("already encoded string - double encoding", "http://url/api"),
Map.of("encoded", "%20"),
"http://url/api?encoded=%2520"
),
Arguments.of(
Named.named("URL with port number", "http://localhost:8080/api/v1/data"),
Map.of("key", "value"),
"http://localhost:8080/api/v1/data?key=value"
),
Arguments.of(
Named.named("URL with userinfo (auth)", "http://user:password@hostname/path"),
Map.of("secure", "true"),
"http://user:password@hostname/path?secure=true"
),
Arguments.of(
Named.named("IPv4 address in URL", "http://192.168.1.100:9090/endpoint"),
Map.of("ip", "test"),
"http://192.168.1.100:9090/endpoint?ip=test"
),
Arguments.of(
Named.named("IPv6 address in URL", "http://[::1]:8080/api"),
Map.of("ipv6", "true"),
"http://[::1]:8080/api?ipv6=true"
),
Arguments.of(
Named.named("HTTPS protocol", "https://secure.example.com/api"),
Map.of("token", "abc123"),
"https://secure.example.com/api?token=abc123"
),
Arguments.of(
Named.named("pipe (|) in value", "http://url/api"),
Map.of("filter", "a|b|c"),
"http://url/api?filter=a%7Cb%7Cc"
),
Arguments.of(
Named.named("caret (^) and backtick (`) in value", "http://url/api"),
Map.of("special", "a^b`c"),
"http://url/api?special=a%5Eb%60c"
),
Arguments.of(
Named.named("tab character in value", "http://url/api"),
Map.of("data", "col1\tcol2\tcol3"),
"http://url/api?data=col1%09col2%09col3"
),
Arguments.of(
Named.named("CRLF in value", "http://url/api"),
Map.of("text", "line1\r\nline2"),
"http://url/api?text=line1%0D%0Aline2"
),
Arguments.of(
Named.named("square brackets in value", "http://url/api"),
Map.of("array", "[1,2,3]"),
"http://url/api?array=%5B1%2C2%2C3%5D"
),
Arguments.of(
Named.named("single and double quotes in value", "http://url/api"),
Map.of("quoted", "He said \"Hello\" and 'Hi'"),
"http://url/api?quoted=He%20said%20%22Hello%22%20and%20%27Hi%27"
),
Arguments.of(
Named.named("angle brackets in value (XSS test)", "http://url/api"),
Map.of("html", "<script>alert('xss')</script>"),
"http://url/api?html=%3Cscript%3Ealert%28%27xss%27%29%3C%2Fscript%3E"
),
Arguments.of(
Named.named("backslash in value", "http://url/api"),
Map.of("path", "C:\\Users\\test"),
"http://url/api?path=C%3A%5CUsers%5Ctest"
),
Arguments.of(
Named.named("at sign (@) in value", "http://url/api"),
Map.of("contact", "user@domain.com"),
"http://url/api?contact=user%40domain.com"
),
Arguments.of(
Named.named("URL as value", "http://url/api"),
Map.of("redirect", "https://example.com/path?foo=bar"),
"http://url/api?redirect=https%3A%2F%2Fexample.com%2Fpath%3Ffoo%3Dbar"
),
Arguments.of(
Named.named("empty map (not null)", "http://url/api"),
Collections.emptyMap(),
"http://url/api"
),
Arguments.of(
Named.named("underscore and hyphen in key/value", "http://url/api"),
Map.of("my-key_name", "my-value_data"),
"http://url/api?my-key_name=my-value_data"
),
Arguments.of(
Named.named("dots in key and value", "http://url/api"),
Map.of("version.major", "1.2.3"),
"http://url/api?version.major=1.2.3"
),
Arguments.of(
Named.named("multiple reserved chars combined", "http://url/api"),
Map.of("complex", "a=1&b=2#section?query"),
"http://url/api?complex=a%3D1%26b%3D2%23section%3Fquery"
),
Arguments.of(
Named.named("trailing slash in URL", "http://url/api/"),
Map.of("param", "value"),
"http://url/api/?param=value"
),
Arguments.of(
Named.named("long value (1000 chars)", "http://url/api"),
Map.of("long", "a".repeat(1000)),
"http://url/api?long=" + "a".repeat(1000)
),
Arguments.of(
Named.named("Chinese characters in value", "http://url/api"),
Map.of("greeting", "你好世界"),
"http://url/api?greeting=%E4%BD%A0%E5%A5%BD%E4%B8%96%E7%95%8C"
),
Arguments.of(
Named.named("Japanese characters in value", "http://url/api"),
Map.of("text", "こんにちは"),
"http://url/api?text=%E3%81%93%E3%82%93%E3%81%AB%E3%81%A1%E3%81%AF"
),
Arguments.of(
Named.named("existing query params in URL", "http://url/api?existing=param"),
Map.of("new", "value"),
"http://url/api?existing=param&new=value"
),
Arguments.of(
Named.named("existing query and fragment in URL", "http://url/api?existing=param#section"),
Map.of("additional", "data"),
"http://url/api?existing=param&additional=data#section"
),
Arguments.of(
Named.named("null byte in value", "http://url/api"),
Map.of("data", "before\u0000after"),
"http://url/api?data=before%00after"
),
Arguments.of(
Named.named("form feed and vertical tab in value", "http://url/api"),
Map.of("whitespace", "a\fb\u000Bc"),
"http://url/api?whitespace=a%0Cb%0Bc"
),
Arguments.of(
Named.named("tilde (unreserved, not encoded)", "http://url/api"),
Map.of("pattern", "~user"),
"http://url/api?pattern=~user"
),
Arguments.of(
Named.named("asterisk in value", "http://url/api"),
Map.of("wildcard", "*.txt"),
"http://url/api?wildcard=%2A.txt"
),
Arguments.of(
Named.named("curly braces in key", "http://url/api"),
Map.of("{key}", "value"),
"http://url/api?%7Bkey%7D=value"
),
Arguments.of(
Named.named("curly braces in value", "http://url/api"),
Map.of("data", "{value}"),
"http://url/api?data=%7Bvalue%7D"
),
Arguments.of(
Named.named("curly braces in both key and value", "http://url/api"),
Map.of("{param}", "{data}"),
"http://url/api?%7Bparam%7D=%7Bdata%7D"
)
); );
} }

114
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/rest/TbRestApiCallNodeTest.java

@ -25,6 +25,7 @@ import org.apache.http.impl.bootstrap.ServerBootstrap;
import org.apache.http.protocol.HttpContext; 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.Test; import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith; import org.junit.jupiter.api.extension.ExtendWith;
import org.junit.jupiter.params.provider.Arguments; import org.junit.jupiter.params.provider.Arguments;
@ -43,24 +44,34 @@ import org.thingsboard.server.common.data.id.EntityId;
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.msg.TbMsgType; import org.thingsboard.server.common.data.msg.TbMsgType;
import org.thingsboard.server.common.data.rule.RuleNode;
import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgDataType; import org.thingsboard.server.common.msg.TbMsgDataType;
import org.thingsboard.server.common.msg.TbMsgMetaData; import org.thingsboard.server.common.msg.TbMsgMetaData;
import org.thingsboard.server.dao.exception.DataValidationException;
import java.io.IOException; import java.io.IOException;
import java.util.Collections; import java.util.Collections;
import java.util.HashMap;
import java.util.Map;
import java.util.concurrent.CountDownLatch; import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
import java.util.stream.Stream; import java.util.stream.Stream;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotSame; import static org.junit.jupiter.api.Assertions.assertNotSame;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertTrue; import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.Mockito.lenient;
import static org.mockito.Mockito.verify; 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;
@ -94,6 +105,14 @@ public class TbRestApiCallNodeTest extends AbstractRuleNodeUpgradeTest {
} }
} }
@BeforeEach
public void setup() {
ruleNode = new RuleNode();
ruleNode.setId(ruleNodeId);
ruleNode.setName("Test REST API call node");
lenient().when(ctx.getSelf()).thenReturn(ruleNode);
}
@AfterEach @AfterEach
public void teardown() { public void teardown() {
if (server != null) { if (server != null) {
@ -101,6 +120,101 @@ public class TbRestApiCallNodeTest extends AbstractRuleNodeUpgradeTest {
} }
} }
@Test
public void shouldNotAllowNullQueryParamValues() {
// GIVEN
var config = new TbRestApiCallNodeConfiguration().defaultConfiguration();
config.setUseNewEncoding(true);
var params = new HashMap<String, String>();
params.put("key", null);
config.setQueryParams(params);
// WHEN-THEN
assertThatThrownBy(() -> new TbRestApiCallNode().init(ctx, new TbNodeConfiguration(JacksonUtil.valueToTree(config))))
.isInstanceOf(TbNodeException.class)
.matches(e -> ((TbNodeException) e).isUnrecoverable())
.rootCause()
.isInstanceOf(DataValidationException.class)
.hasMessageContaining("query parameter values must be non-null");
}
@Test
public void shouldAllowOnlyNullQueryParamsIfOldEncodingIsUsed() {
// GIVEN
var config = new TbRestApiCallNodeConfiguration().defaultConfiguration();
config.setUseNewEncoding(false);
config.setQueryParams(Map.of("key", "value"));
// WHEN-THEN
assertThatThrownBy(() -> new TbRestApiCallNode().init(ctx, new TbNodeConfiguration(JacksonUtil.valueToTree(config))))
.isInstanceOf(TbNodeException.class)
.hasRootCauseInstanceOf(DataValidationException.class)
.hasRootCauseMessage("'" + ruleNode.getName() + "' node configuration is invalid: query parameters must be null if old encoding is used")
.matches(e -> ((TbNodeException) e).isUnrecoverable());
}
@Test
public void shouldNotAllowNullQueryParamsIfNewEncodingIsUsed() {
// GIVEN
var config = new TbRestApiCallNodeConfiguration().defaultConfiguration();
config.setUseNewEncoding(true);
config.setQueryParams(null);
// WHEN-THEN
assertThatThrownBy(() -> new TbRestApiCallNode().init(ctx, new TbNodeConfiguration(JacksonUtil.valueToTree(config))))
.isInstanceOf(TbNodeException.class)
.hasRootCauseInstanceOf(DataValidationException.class)
.hasRootCauseMessage("'" + ruleNode.getName() + "' node configuration is invalid: query parameters must be non-null if new encoding is used")
.matches(e -> ((TbNodeException) e).isUnrecoverable());
}
@Test
public void shouldUseNewEncodingForNewNodesByDefault() {
// GIVEN-WHEN
var defaultConfig = new TbRestApiCallNodeConfiguration().defaultConfiguration();
// THEN
assertTrue(defaultConfig.isUseNewEncoding());
assertEquals(Collections.emptyMap(), defaultConfig.getQueryParams());
}
@Test
public void shouldDefaultToOldEncodingForLegacyConfigs() {
// GIVEN
String configJson = """
{
"restEndpointUrlPattern": "http://url?param=value",
"requestMethod": "GET",
"useSimpleClientHttpFactory": false,
"parseToPlainText": false,
"ignoreRequestBody": false,
"enableProxy": false,
"useSystemProxyProperties": false,
"proxyScheme": null,
"proxyHost": null,
"proxyPort": 0,
"proxyUser": null,
"proxyPassword": null,
"readTimeoutMs": 0,
"maxParallelRequestsCount": 0,
"headers": {
"Content-Type": "application/json",
"X-Authorization": "Bearer eyJhbGciOi...."
},
"credentials": {
"type": "anonymous"
},
"maxInMemoryBufferSizeInKb": 256
}""";
// WHEN
var config = JacksonUtil.fromString(configJson, TbRestApiCallNodeConfiguration.class);
// THEN
assertFalse(config.isUseNewEncoding());
assertNull(config.getQueryParams());
}
@Test @Test
public void deleteRequestWithoutBody() throws IOException, InterruptedException { public void deleteRequestWithoutBody() throws IOException, InterruptedException {
final CountDownLatch latch = new CountDownLatch(1); final CountDownLatch latch = new CountDownLatch(1);

Loading…
Cancel
Save