diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbHttpClient.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbHttpClient.java index 96ec648941..c21914577f 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbHttpClient.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbHttpClient.java @@ -80,6 +80,8 @@ public class TbHttpClient { public static final String PROXY_USER = "tb.proxy.user"; public static final String PROXY_PASSWORD = "tb.proxy.password"; + public static final String MAX_IN_MEMORY_BUFFER_SIZE_IN_KB = "tb.http.maxInMemoryBufferSizeInKb"; + private final TbRestApiCallNodeConfiguration config; private EventLoopGroup eventLoopGroup; @@ -129,15 +131,34 @@ public class TbHttpClient { httpClient = httpClient.secure(t -> t.sslContext(sslContext)); } + validateMaxInMemoryBufferSize(config); + this.webClient = WebClient.builder() .clientConnector(new ReactorClientHttpConnector(httpClient)) .defaultHeader(HttpHeaders.CONNECTION, "close") //In previous realization this header was present! (Added for hotfix "Connection reset") + .codecs(configurer -> configurer.defaultCodecs().maxInMemorySize( + (config.getMaxInMemoryBufferSizeInKb() > 0 ? config.getMaxInMemoryBufferSizeInKb() : 256) * 1024)) .build(); } catch (SSLException e) { throw new TbNodeException(e); } } + private void validateMaxInMemoryBufferSize(TbRestApiCallNodeConfiguration config) throws TbNodeException { + int systemMaxInMemoryBufferSizeInKb = 25000; + try { + Properties properties = System.getProperties(); + if (properties.containsKey(MAX_IN_MEMORY_BUFFER_SIZE_IN_KB)) { + systemMaxInMemoryBufferSizeInKb = Integer.parseInt(properties.getProperty(MAX_IN_MEMORY_BUFFER_SIZE_IN_KB)); + } + } catch (Exception ignored) {} + if (config.getMaxInMemoryBufferSizeInKb() > systemMaxInMemoryBufferSizeInKb) { + throw new TbNodeException("The configured maximum in-memory buffer size (in KB) exceeds the system limit for this parameter.\n" + + "The system limit is " + systemMaxInMemoryBufferSizeInKb + " KB.\n" + + "Please use the system variable '" + MAX_IN_MEMORY_BUFFER_SIZE_IN_KB + "' to override the system limit."); + } + } + EventLoopGroup getSharedOrCreateEventLoopGroup(EventLoopGroup eventLoopGroupShared) { if (eventLoopGroupShared != null) { return eventLoopGroupShared; @@ -207,13 +228,23 @@ public class TbHttpClient { semaphore.release(); } - onFailure.accept(processException(msg, throwable), throwable); + onFailure.accept(processException(msg, throwable), processThrowable(throwable)); }); } catch (InterruptedException e) { log.warn("Timeout during waiting for reply!", e); } } + private Throwable processThrowable(Throwable origin) { + if (origin instanceof WebClientResponseException restClientResponseException + && restClientResponseException.getStatusCode().is2xxSuccessful()) { + // return cause instead of original exception in case 2xx status code + // this will provide meaningful error message to the user + return new RuntimeException(restClientResponseException.getCause()); + } + return origin; + } + public URI buildEncodedUri(String endpointUrl) { if (endpointUrl == null) { throw new RuntimeException("Url string cannot be null!"); @@ -336,8 +367,8 @@ public class TbHttpClient { String hostname = properties.getProperty(hostProperty); int port = Integer.parseInt(properties.getProperty(portProperty)); - checkProxyHost(config.getProxyHost()); - checkProxyPort(config.getProxyPort()); + checkProxyHost(hostname); + checkProxyPort(port); var proxy = option .type(ProxyProvider.Proxy.HTTP) @@ -362,8 +393,8 @@ public class TbHttpClient { ProxyProvider.Proxy type = SOCKS_VERSION_5.equals(version) ? ProxyProvider.Proxy.SOCKS5 : ProxyProvider.Proxy.SOCKS4; int port = Integer.parseInt(properties.getProperty(SOCKS_PROXY_PORT)); - checkProxyHost(config.getProxyHost()); - checkProxyPort(config.getProxyPort()); + checkProxyHost(hostname); + checkProxyPort(port); ProxyProvider.Builder proxy = option .type(type) diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbRestApiCallNodeConfiguration.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbRestApiCallNodeConfiguration.java index 7d2ff7167d..1825bd1ae7 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbRestApiCallNodeConfiguration.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbRestApiCallNodeConfiguration.java @@ -46,6 +46,7 @@ public class TbRestApiCallNodeConfiguration implements NodeConfiguration