From 7a140b25181488e9855fe807c6641afe46f788e1 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Tue, 2 Jan 2024 20:13:35 +0100 Subject: [PATCH 01/21] haproxy limits, blocklist, trustlist --- docker/haproxy/config/blocklist.txt | 3 + docker/haproxy/config/haproxy.cfg | 102 ++++++++++++++++++++++++++-- docker/haproxy/config/trustlist.txt | 12 ++++ 3 files changed, 111 insertions(+), 6 deletions(-) create mode 100644 docker/haproxy/config/blocklist.txt create mode 100644 docker/haproxy/config/trustlist.txt diff --git a/docker/haproxy/config/blocklist.txt b/docker/haproxy/config/blocklist.txt new file mode 100644 index 0000000000..ff9429857f --- /dev/null +++ b/docker/haproxy/config/blocklist.txt @@ -0,0 +1,3 @@ +# Blocked subnets and IPs. Use CIDR or IP by one per line +5.136.0.0/13 +217.199.254.1 diff --git a/docker/haproxy/config/haproxy.cfg b/docker/haproxy/config/haproxy.cfg index aa752502bf..757efe692f 100644 --- a/docker/haproxy/config/haproxy.cfg +++ b/docker/haproxy/config/haproxy.cfg @@ -41,6 +41,15 @@ listen stats listen mqtt-in bind *:${MQTT_PORT} mode tcp + + stick-table type ip size 60k expire 60s store conn_cur + + acl trustlist src -f /config/trustlist.txt + acl blocklist src -f /config/blocklist.txt + tcp-request connection accept if trustlist + tcp-request connection reject if blocklist or { src_conn_cur ge 50 } + tcp-request connection track-sc1 src + option clitcpka # For TCP keep-alive timeout client 3h timeout server 3h @@ -52,6 +61,15 @@ listen mqtt-in listen edges-rpc-in bind *:${EDGES_RPC_PORT} mode tcp + + stick-table type ip size 60k expire 60s store conn_cur + + acl trustlist src -f /config/trustlist.txt + acl blocklist src -f /config/blocklist.txt + tcp-request connection accept if trustlist + tcp-request connection reject if blocklist or { src_conn_cur ge 5 } + tcp-request connection track-sc1 src + option clitcpka # For TCP keep-alive timeout client 3h timeout server 3h @@ -63,18 +81,28 @@ listen edges-rpc-in frontend http-in bind *:${HTTP_PORT} alpn h2,http/1.1 + stick-table type ip size 60k expire 60s store conn_cur + + acl trustlist src -f /config/trustlist.txt + acl blocklist src -f /config/blocklist.txt + tcp-request connection accept if trustlist + tcp-request connection reject if blocklist or { src_conn_cur ge 50 } + tcp-request connection track-sc1 src + option forwardfor http-request add-header "X-Forwarded-Proto" "http" acl transport_http_acl path_beg /api/v1/ acl letsencrypt_http_acl path_beg /.well-known/acme-challenge/ + acl tb_images_api_acl path_beg /api/images/ acl tb_api_acl path_beg /api/ /swagger /webjars /v2/ /v3/ /static/rulenode/ /oauth2/ /login/oauth2/ /static/widgets/ redirect scheme https if !letsencrypt_http_acl !transport_http_acl { env(FORCE_HTTPS_REDIRECT) -m str true } use_backend letsencrypt_http if letsencrypt_http_acl use_backend tb-http-backend if transport_http_acl + use_backend tb-images-api-backend if tb_images_api_acl use_backend tb-api-backend if tb_api_acl default_backend tb-web-backend @@ -82,14 +110,24 @@ frontend http-in frontend https_in bind *:${HTTPS_PORT} ssl crt /usr/local/etc/haproxy/default.pem crt /usr/local/etc/haproxy/certs.d ciphers ECDHE-RSA-AES256-SHA:RC4-SHA:RC4:HIGH:!MD5:!aNULL:!EDH:!AESGCM alpn h2,http/1.1 + stick-table type ip size 60k expire 60s store conn_cur + + acl trustlist src -f /config/trustlist.txt + acl blocklist src -f /config/blocklist.txt + tcp-request connection accept if trustlist + tcp-request connection reject if blocklist or { src_conn_cur ge 50 } + tcp-request connection track-sc1 src + option forwardfor http-request add-header "X-Forwarded-Proto" "https" acl transport_http_acl path_beg /api/v1/ + acl tb_images_api_acl path_beg /api/images/ acl tb_api_acl path_beg /api/ /swagger /webjars /v2/ /v3/ /static/rulenode/ /oauth2/ /login/oauth2/ /static/widgets/ use_backend tb-http-backend if transport_http_acl + use_backend tb-images-api-backend if tb_images_api_acl use_backend tb-api-backend if tb_api_acl default_backend tb-web-backend @@ -98,24 +136,76 @@ backend letsencrypt_http server letsencrypt_http_srv 127.0.0.1:8080 backend tb-web-backend + timeout queue 60s balance leastconn option tcp-check option log-health-checks - server tbWeb1 tb-web-ui1:8080 check inter 5s resolvers docker_resolver resolve-prefer ipv4 - server tbWeb2 tb-web-ui2:8080 check inter 5s resolvers docker_resolver resolve-prefer ipv4 + server tbWeb1 tb-web-ui1:8080 check inter 5s resolvers docker_resolver resolve-prefer ipv4 maxconn 50 + server tbWeb2 tb-web-ui2:8080 check inter 5s resolvers docker_resolver resolve-prefer ipv4 maxconn 50 http-request set-header X-Forwarded-Port %[dst_port] backend tb-http-backend + timeout queue 60s balance leastconn option tcp-check option log-health-checks - server tbHttp1 tb-http-transport1:8081 check inter 5s resolvers docker_resolver resolve-prefer ipv4 - server tbHttp2 tb-http-transport2:8081 check inter 5s resolvers docker_resolver resolve-prefer ipv4 + server tbHttp1 tb-http-transport1:8081 check inter 5s resolvers docker_resolver resolve-prefer ipv4 maxconn 50 + server tbHttp2 tb-http-transport2:8081 check inter 5s resolvers docker_resolver resolve-prefer ipv4 maxconn 50 + +# Dummy backends for a stick-table purpose only. +# There is only one stick-table per proxy. At the moment of writing this doc, +# it does not seem useful to have multiple tables per proxy. If this happens +# to be required, simply create a dummy backend with a stick-table in it and +# reference it +# https://www.haproxy.com/documentation/haproxy-configuration-manual/latest/#stick-table +backend st_src_rate10s + stick-table type ip size 60k expire 10s store http_req_rate(10s) + +backend st_src_rate1m + stick-table type ip size 60k expire 1m store http_req_rate(1m) backend tb-api-backend + timeout queue 60s balance source option tcp-check option log-health-checks - server tbApi1 tb-core1:8080 check inter 5s resolvers docker_resolver resolve-prefer ipv4 - server tbApi2 tb-core2:8080 check inter 5s resolvers docker_resolver resolve-prefer ipv4 + + http-request track-sc0 src table st_src_rate10s + http-request track-sc1 src table st_src_rate1m + + acl trustlist src -f /config/trustlist.txt + http-request deny deny_status 429 if { sc_http_req_rate(0) gt 100 } !trustlist + http-request deny deny_status 429 if { sc_http_req_rate(1) gt 300 } !trustlist + + http-request set-header X-Forwarded-Port %[dst_port] + server tbApi1 tb-core1:8080 check inter 5s resolvers docker_resolver resolve-prefer ipv4 maxconn 50 + server tbApi2 tb-core2:8080 check inter 5s resolvers docker_resolver resolve-prefer ipv4 maxconn 50 + +# Dummy backends for a stick-table purpose only. +# There is only one stick-table per proxy. At the moment of writing this doc, +# it does not seem useful to have multiple tables per proxy. If this happens +# to be required, simply create a dummy backend with a stick-table in it and +# reference it +# https://www.haproxy.com/documentation/haproxy-configuration-manual/latest/#stick-table +backend st_images_src_rate10s + stick-table type ip size 60k expire 10s store http_req_rate(10s) + +backend st_images_src_rate1m + stick-table type ip size 60k expire 1m store http_req_rate(1m) + +backend tb-images-api-backend + timeout queue 60s + balance source + option tcp-check + option log-health-checks + + http-request track-sc0 src table st_images_src_rate10s + http-request track-sc1 src table st_images_src_rate1m + + acl trustlist src -f /config/trustlist.txt + http-request deny deny_status 429 if { sc_http_req_rate(0) gt 1000 } !trustlist + http-request deny deny_status 429 if { sc_http_req_rate(1) gt 3000 } !trustlist + http-request set-header X-Forwarded-Port %[dst_port] + server tbImagesApi1 tb-core1:8080 check inter 10s resolvers docker_resolver resolve-prefer ipv4 maxconn 50 + server tbImagesApi2 tb-core2:8080 check inter 10s resolvers docker_resolver resolve-prefer ipv4 maxconn 50 diff --git a/docker/haproxy/config/trustlist.txt b/docker/haproxy/config/trustlist.txt new file mode 100644 index 0000000000..93716b46d5 --- /dev/null +++ b/docker/haproxy/config/trustlist.txt @@ -0,0 +1,12 @@ +# Trusted list is intended to do not apply any limitations for trustees +# +# Private subnet example +# 10.0.0.0/8 +# Docker-compose subnet +172.16.0.0/12 +# Local network subnet +192.168.0.0/16 +# Allow loopback interface +127.0.0.1 +::1 +# Allow trusted IPs or CIDRs below From ee1209b19f4e842aff42a07b81d9c579574d2f32 Mon Sep 17 00:00:00 2001 From: Andrii Landiak Date: Fri, 22 Mar 2024 12:01:36 +0200 Subject: [PATCH 02/21] Fix KvProtoUtils order for matching KeyValueType and DataType --- .../device/DeviceActorMessageProcessor.java | 5 +- .../queue/DefaultTbCoreConsumerService.java | 8 +- .../server/common/data/EntityType.java | 3 +- .../server/common/data/kv/DataType.java | 15 +- .../server/common/util/KvProtoUtil.java | 128 +++++++---------- .../server/common/util/KvProtoUtilTest.java | 131 ++++++++++++++++++ 6 files changed, 202 insertions(+), 88 deletions(-) create mode 100644 common/proto/src/test/java/org/thingsboard/server/common/util/KvProtoUtilTest.java diff --git a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java index d6b5da4ddf..d5fc8055d3 100644 --- a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java @@ -21,6 +21,7 @@ import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.MoreExecutors; +import jakarta.annotation.Nullable; import lombok.extern.slf4j.Slf4j; import org.apache.commons.collections.CollectionUtils; import org.thingsboard.common.util.JacksonUtil; @@ -68,7 +69,6 @@ import org.thingsboard.server.common.msg.rule.engine.DeviceEdgeUpdateMsg; import org.thingsboard.server.common.msg.rule.engine.DeviceNameOrTypeUpdateMsg; import org.thingsboard.server.common.msg.timeout.DeviceActorServerSideRpcTimeoutMsg; import org.thingsboard.server.common.util.KvProtoUtil; -import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.AttributeUpdateNotificationMsg; import org.thingsboard.server.gen.transport.TransportProtos.ClaimDeviceMsg; import org.thingsboard.server.gen.transport.TransportProtos.DeviceSessionsCacheEntry; @@ -97,7 +97,6 @@ import org.thingsboard.server.service.rpc.RpcSubmitStrategy; import org.thingsboard.server.service.state.DefaultDeviceStateService; import org.thingsboard.server.service.transport.msg.TransportToDeviceActorMsgWrapper; -import jakarta.annotation.Nullable; import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; @@ -605,7 +604,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso if (DataConstants.SHARED_SCOPE.equals(msg.getScope())) { List attributes = new ArrayList<>(msg.getValues()); if (attributes.size() > 0) { - List sharedUpdated = msg.getValues().stream().map(t -> KvProtoUtil.toTsKvProto(t.getLastUpdateTs(), t)) + List sharedUpdated = msg.getValues().stream().map(t -> KvProtoUtil.toProto(t.getLastUpdateTs(), t)) .collect(Collectors.toList()); if (!sharedUpdated.isEmpty()) { notification.addAllSharedUpdated(sharedUpdated); diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java index fd08b9da4d..fbbb65d4ef 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java @@ -18,6 +18,8 @@ package org.thingsboard.server.service.queue; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListeningExecutorService; import com.google.common.util.concurrent.MoreExecutors; +import jakarta.annotation.PostConstruct; +import jakarta.annotation.PreDestroy; import lombok.Getter; import lombok.Setter; import lombok.extern.slf4j.Slf4j; @@ -50,9 +52,9 @@ import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.common.msg.rpc.FromDeviceRpcResponse; import org.thingsboard.server.common.msg.rpc.ToDeviceRpcRequestActorMsg; import org.thingsboard.server.common.stats.StatsFactory; +import org.thingsboard.server.common.util.KvProtoUtil; import org.thingsboard.server.common.util.ProtoUtils; import org.thingsboard.server.dao.resource.ImageCacheKey; -import org.thingsboard.server.common.util.KvProtoUtil; import org.thingsboard.server.dao.tenant.TbTenantProfileCache; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.DeviceStateServiceMsgProto; @@ -101,8 +103,6 @@ import org.thingsboard.server.service.transport.msg.TransportToDeviceActorMsgWra import org.thingsboard.server.service.ws.notification.sub.NotificationRequestUpdate; import org.thingsboard.server.service.ws.notification.sub.NotificationUpdate; -import jakarta.annotation.PostConstruct; -import jakarta.annotation.PreDestroy; import java.util.List; import java.util.Optional; import java.util.UUID; @@ -583,7 +583,7 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService dataTypeByProtoNumber[dataType.getProtoNumber()] = dataType); + } + + public static List toAttributeKvList(List dataList) { + List result = new ArrayList<>(dataList.size()); + dataList.forEach(proto -> result.add(new BaseAttributeKvEntry(fromProto(proto.getKv()), proto.getTs()))); + return result; + } + public static List attrToTsKvProtos(List result) { List clientAttributes; if (result == null || result.isEmpty()) { @@ -41,115 +56,70 @@ public class KvProtoUtil { } else { clientAttributes = new ArrayList<>(result.size()); for (AttributeKvEntry attrEntry : result) { - clientAttributes.add(toTsKvProto(attrEntry.getLastUpdateTs(), attrEntry)); + clientAttributes.add(toProto(attrEntry.getLastUpdateTs(), attrEntry)); } } return clientAttributes; } - - public static List tsToTsKvProtos(List result) { + public static List toProtoList(List result) { List ts; if (result == null || result.isEmpty()) { ts = Collections.emptyList(); } else { ts = new ArrayList<>(result.size()); for (TsKvEntry attrEntry : result) { - ts.add(toTsKvProto(attrEntry.getTs(), attrEntry)); + ts.add(toProto(attrEntry.getTs(), attrEntry)); } } return ts; } - public static TransportProtos.TsKvProto toTsKvProto(long ts, KvEntry kvEntry) { + public static List fromProtoList(List dataList) { + List result = new ArrayList<>(dataList.size()); + dataList.forEach(proto -> result.add(new BasicTsKvEntry(proto.getTs(), fromProto(proto.getKv())))); + return result; + } + + public static TransportProtos.TsKvProto toProto(long ts, KvEntry kvEntry) { return TransportProtos.TsKvProto.newBuilder().setTs(ts) - .setKv(KvProtoUtil.toKeyValueProto(kvEntry)).build(); + .setKv(KvProtoUtil.toProto(kvEntry)).build(); + } + + public static TsKvEntry fromProto(TransportProtos.TsKvProto proto) { + return new BasicTsKvEntry(proto.getTs(), fromProto(proto.getKv())); } - public static TransportProtos.KeyValueProto toKeyValueProto(KvEntry kvEntry) { + public static TransportProtos.KeyValueProto toProto(KvEntry kvEntry) { TransportProtos.KeyValueProto.Builder builder = TransportProtos.KeyValueProto.newBuilder(); builder.setKey(kvEntry.getKey()); + builder.setType(toProto(kvEntry.getDataType())); switch (kvEntry.getDataType()) { - case BOOLEAN: - builder.setType(TransportProtos.KeyValueType.BOOLEAN_V); - builder.setBoolV(kvEntry.getBooleanValue().get()); - break; - case DOUBLE: - builder.setType(TransportProtos.KeyValueType.DOUBLE_V); - builder.setDoubleV(kvEntry.getDoubleValue().get()); - break; - case LONG: - builder.setType(TransportProtos.KeyValueType.LONG_V); - builder.setLongV(kvEntry.getLongValue().get()); - break; - case STRING: - builder.setType(TransportProtos.KeyValueType.STRING_V); - builder.setStringV(kvEntry.getStrValue().get()); - break; - case JSON: - builder.setType(TransportProtos.KeyValueType.JSON_V); - builder.setJsonV(kvEntry.getJsonValue().get()); - break; + case BOOLEAN -> kvEntry.getBooleanValue().ifPresent(builder::setBoolV); + case LONG -> kvEntry.getLongValue().ifPresent(builder::setLongV); + case DOUBLE -> kvEntry.getDoubleValue().ifPresent(builder::setDoubleV); + case JSON -> kvEntry.getJsonValue().ifPresent(builder::setJsonV); + case STRING -> kvEntry.getStrValue().ifPresent(builder::setStringV); } return builder.build(); } - public static TransportProtos.TsKvProto.Builder toKeyValueProto(long ts, KvEntry attr) { - TransportProtos.KeyValueProto.Builder dataBuilder = TransportProtos.KeyValueProto.newBuilder(); - dataBuilder.setKey(attr.getKey()); - dataBuilder.setType(TransportProtos.KeyValueType.forNumber(attr.getDataType().ordinal())); - switch (attr.getDataType()) { - case BOOLEAN: - attr.getBooleanValue().ifPresent(dataBuilder::setBoolV); - break; - case LONG: - attr.getLongValue().ifPresent(dataBuilder::setLongV); - break; - case DOUBLE: - attr.getDoubleValue().ifPresent(dataBuilder::setDoubleV); - break; - case JSON: - attr.getJsonValue().ifPresent(dataBuilder::setJsonV); - break; - case STRING: - attr.getStrValue().ifPresent(dataBuilder::setStringV); - break; - } - return TransportProtos.TsKvProto.newBuilder().setTs(ts).setKv(dataBuilder); + public static KvEntry fromProto(TransportProtos.KeyValueProto proto) { + return switch (fromProto(proto.getType())) { + case BOOLEAN -> new BooleanDataEntry(proto.getKey(), proto.getBoolV()); + case LONG -> new LongDataEntry(proto.getKey(), proto.getLongV()); + case DOUBLE -> new DoubleDataEntry(proto.getKey(), proto.getDoubleV()); + case STRING -> new StringDataEntry(proto.getKey(), proto.getStringV()); + case JSON -> new JsonDataEntry(proto.getKey(), proto.getJsonV()); + }; } - public static List toTsKvEntityList(List dataList) { - List result = new ArrayList<>(dataList.size()); - dataList.forEach(proto -> result.add(new BasicTsKvEntry(proto.getTs(), getKvEntry(proto.getKv())))); - return result; + public static TransportProtos.KeyValueType toProto(DataType dataType) { + return TransportProtos.KeyValueType.forNumber(dataType.getProtoNumber()); } - public static List toAttributeKvList(List dataList) { - List result = new ArrayList<>(dataList.size()); - dataList.forEach(proto -> result.add(new BaseAttributeKvEntry(getKvEntry(proto.getKv()), proto.getTs()))); - return result; + public static DataType fromProto(TransportProtos.KeyValueType keyValueType) { + return dataTypeByProtoNumber[keyValueType.getNumber()]; } - private static KvEntry getKvEntry(TransportProtos.KeyValueProto proto) { - KvEntry entry = null; - DataType type = DataType.values()[proto.getType().getNumber()]; - switch (type) { - case BOOLEAN: - entry = new BooleanDataEntry(proto.getKey(), proto.getBoolV()); - break; - case LONG: - entry = new LongDataEntry(proto.getKey(), proto.getLongV()); - break; - case DOUBLE: - entry = new DoubleDataEntry(proto.getKey(), proto.getDoubleV()); - break; - case STRING: - entry = new StringDataEntry(proto.getKey(), proto.getStringV()); - break; - case JSON: - entry = new JsonDataEntry(proto.getKey(), proto.getJsonV()); - break; - } - return entry; - } } diff --git a/common/proto/src/test/java/org/thingsboard/server/common/util/KvProtoUtilTest.java b/common/proto/src/test/java/org/thingsboard/server/common/util/KvProtoUtilTest.java new file mode 100644 index 0000000000..603d4b3aa6 --- /dev/null +++ b/common/proto/src/test/java/org/thingsboard/server/common/util/KvProtoUtilTest.java @@ -0,0 +1,131 @@ +/** + * Copyright © 2016-2024 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.util; + +import org.junit.jupiter.api.Test; +import org.thingsboard.server.common.data.kv.AggTsKvEntry; +import org.thingsboard.server.common.data.kv.AttributeKvEntry; +import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; +import org.thingsboard.server.common.data.kv.BasicTsKvEntry; +import org.thingsboard.server.common.data.kv.BooleanDataEntry; +import org.thingsboard.server.common.data.kv.DataType; +import org.thingsboard.server.common.data.kv.DoubleDataEntry; +import org.thingsboard.server.common.data.kv.JsonDataEntry; +import org.thingsboard.server.common.data.kv.KvEntry; +import org.thingsboard.server.common.data.kv.LongDataEntry; +import org.thingsboard.server.common.data.kv.StringDataEntry; +import org.thingsboard.server.common.data.kv.TsKvEntry; + +import java.util.List; + +import static org.assertj.core.api.Assertions.assertThat; + +class KvProtoUtilTest { + + @Test + void protoDataTypeSerialization() { + for (DataType dataType : DataType.values()) { + assertThat(KvProtoUtil.fromProto(KvProtoUtil.toProto(dataType))).as(dataType.name()).isEqualTo(dataType); + } + } + + @Test + void protoKeyValueProtoSerialization() { + String key = "key"; + KvEntry kvEntry = new BooleanDataEntry(key, true); + assertThat(KvProtoUtil.fromProto(KvProtoUtil.toProto(kvEntry))).as("deserialized").isEqualTo(kvEntry); + + kvEntry = new LongDataEntry(key, 23L); + assertThat(KvProtoUtil.fromProto(KvProtoUtil.toProto(kvEntry))).as("deserialized").isEqualTo(kvEntry); + + kvEntry = new DoubleDataEntry(key, 23.0); + assertThat(KvProtoUtil.fromProto(KvProtoUtil.toProto(kvEntry))).as("deserialized").isEqualTo(kvEntry); + + kvEntry = new StringDataEntry(key, "stringValue"); + assertThat(KvProtoUtil.fromProto(KvProtoUtil.toProto(kvEntry))).as("deserialized").isEqualTo(kvEntry); + + kvEntry = new JsonDataEntry(key, "jsonValue"); + assertThat(KvProtoUtil.fromProto(KvProtoUtil.toProto(kvEntry))).as("deserialized").isEqualTo(kvEntry); + } + + @Test + void protoTsKvEntrySerialization() { + String key = "key"; + long ts = System.currentTimeMillis(); + KvEntry kvEntry = new BasicTsKvEntry(ts, new BooleanDataEntry(key, true)); + assertThat(KvProtoUtil.fromProto(KvProtoUtil.toProto(ts, kvEntry))).as("deserialized").isEqualTo(kvEntry); + + kvEntry = new BasicTsKvEntry(ts, new LongDataEntry(key, 23L)); + assertThat(KvProtoUtil.fromProto(KvProtoUtil.toProto(ts, kvEntry))).as("deserialized").isEqualTo(kvEntry); + + kvEntry = new BasicTsKvEntry(ts, new DoubleDataEntry(key, 23.0)); + assertThat(KvProtoUtil.fromProto(KvProtoUtil.toProto(ts, kvEntry))).as("deserialized").isEqualTo(kvEntry); + + kvEntry = new BasicTsKvEntry(ts, new StringDataEntry(key, "stringValue")); + assertThat(KvProtoUtil.fromProto(KvProtoUtil.toProto(ts, kvEntry))).as("deserialized").isEqualTo(kvEntry); + + kvEntry = new BasicTsKvEntry(ts, new JsonDataEntry(key, "jsonValue")); + assertThat(KvProtoUtil.fromProto(KvProtoUtil.toProto(ts, kvEntry))).as("deserialized").isEqualTo(kvEntry); + } + + @Test + void protoListTsKvEntrySerialization() { + String key = "key"; + long ts = System.currentTimeMillis(); + KvEntry booleanDataEntry = new BooleanDataEntry(key, true); + KvEntry longDataEntry = new LongDataEntry(key, 23L); + KvEntry doubleDataEntry = new DoubleDataEntry(key, 23.0); + KvEntry stringDataEntry = new StringDataEntry(key, "stringValue"); + KvEntry jsonDataEntry = new JsonDataEntry(key, "jsonValue"); + List protoList = List.of( + new BasicTsKvEntry(ts, booleanDataEntry), + new BasicTsKvEntry(ts, longDataEntry), + new BasicTsKvEntry(ts, doubleDataEntry), + new BasicTsKvEntry(ts, stringDataEntry), + new BasicTsKvEntry(ts, jsonDataEntry) + ); + assertThat(KvProtoUtil.fromProtoList(KvProtoUtil.toProtoList(protoList))).as("deserialized").isEqualTo(protoList); + + protoList = List.of( + new AggTsKvEntry(ts, booleanDataEntry, 3), + new AggTsKvEntry(ts, longDataEntry, 5), + new AggTsKvEntry(ts, doubleDataEntry, 2), + new AggTsKvEntry(ts, stringDataEntry, 1), + new AggTsKvEntry(ts, jsonDataEntry, 0) + ); + assertThat(KvProtoUtil.fromProtoList(KvProtoUtil.toProtoList(protoList))).as("deserialized").isEqualTo(protoList); + } + + @Test + void protoListAttributeKvSerialization() { + String key = "key"; + long ts = System.currentTimeMillis(); + KvEntry booleanDataEntry = new BooleanDataEntry(key, true); + KvEntry longDataEntry = new LongDataEntry(key, 23L); + KvEntry doubleDataEntry = new DoubleDataEntry(key, 23.0); + KvEntry stringDataEntry = new StringDataEntry(key, "stringValue"); + KvEntry jsonDataEntry = new JsonDataEntry(key, "jsonValue"); + List protoList = List.of( + new BaseAttributeKvEntry(ts, booleanDataEntry), + new BaseAttributeKvEntry(ts, longDataEntry), + new BaseAttributeKvEntry(ts, doubleDataEntry), + new BaseAttributeKvEntry(ts, stringDataEntry), + new BaseAttributeKvEntry(ts, jsonDataEntry) + ); + assertThat(KvProtoUtil.toAttributeKvList(KvProtoUtil.attrToTsKvProtos(protoList))).as("deserialized").isEqualTo(protoList); + } + +} From 7edb18673c9d8d5913e7b00551da0279c4340211 Mon Sep 17 00:00:00 2001 From: Andrii Landiak Date: Mon, 25 Mar 2024 12:47:20 +0200 Subject: [PATCH 03/21] Improve test for KvProtoUtil --- .../server/common/util/KvProtoUtilTest.java | 135 ++++++++---------- 1 file changed, 59 insertions(+), 76 deletions(-) diff --git a/common/proto/src/test/java/org/thingsboard/server/common/util/KvProtoUtilTest.java b/common/proto/src/test/java/org/thingsboard/server/common/util/KvProtoUtilTest.java index 603d4b3aa6..0fad4e48a2 100644 --- a/common/proto/src/test/java/org/thingsboard/server/common/util/KvProtoUtilTest.java +++ b/common/proto/src/test/java/org/thingsboard/server/common/util/KvProtoUtilTest.java @@ -16,6 +16,10 @@ package org.thingsboard.server.common.util; import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.EnumSource; +import org.junit.jupiter.params.provider.MethodSource; +import org.junit.jupiter.params.provider.ValueSource; import org.thingsboard.server.common.data.kv.AggTsKvEntry; import org.thingsboard.server.common.data.kv.AttributeKvEntry; import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; @@ -30,102 +34,81 @@ import org.thingsboard.server.common.data.kv.StringDataEntry; import org.thingsboard.server.common.data.kv.TsKvEntry; import java.util.List; +import java.util.stream.Collectors; +import java.util.stream.Stream; import static org.assertj.core.api.Assertions.assertThat; class KvProtoUtilTest { - @Test - void protoDataTypeSerialization() { - for (DataType dataType : DataType.values()) { - assertThat(KvProtoUtil.fromProto(KvProtoUtil.toProto(dataType))).as(dataType.name()).isEqualTo(dataType); - } - } + private static final long TS = System.currentTimeMillis(); - @Test - void protoKeyValueProtoSerialization() { + private static Stream kvEntryData() { String key = "key"; - KvEntry kvEntry = new BooleanDataEntry(key, true); - assertThat(KvProtoUtil.fromProto(KvProtoUtil.toProto(kvEntry))).as("deserialized").isEqualTo(kvEntry); - - kvEntry = new LongDataEntry(key, 23L); - assertThat(KvProtoUtil.fromProto(KvProtoUtil.toProto(kvEntry))).as("deserialized").isEqualTo(kvEntry); - - kvEntry = new DoubleDataEntry(key, 23.0); - assertThat(KvProtoUtil.fromProto(KvProtoUtil.toProto(kvEntry))).as("deserialized").isEqualTo(kvEntry); - - kvEntry = new StringDataEntry(key, "stringValue"); - assertThat(KvProtoUtil.fromProto(KvProtoUtil.toProto(kvEntry))).as("deserialized").isEqualTo(kvEntry); - - kvEntry = new JsonDataEntry(key, "jsonValue"); - assertThat(KvProtoUtil.fromProto(KvProtoUtil.toProto(kvEntry))).as("deserialized").isEqualTo(kvEntry); + return Stream.of( + new BooleanDataEntry(key, true), + new LongDataEntry(key, 23L), + new DoubleDataEntry(key, 23.0), + new StringDataEntry(key, "stringValue"), + new JsonDataEntry(key, "jsonValue") + ); } - @Test - void protoTsKvEntrySerialization() { - String key = "key"; - long ts = System.currentTimeMillis(); - KvEntry kvEntry = new BasicTsKvEntry(ts, new BooleanDataEntry(key, true)); - assertThat(KvProtoUtil.fromProto(KvProtoUtil.toProto(ts, kvEntry))).as("deserialized").isEqualTo(kvEntry); + private static Stream basicTsKvEntryData() { + return kvEntryData().map(kvEntry -> new BasicTsKvEntry(TS, kvEntry)); + } - kvEntry = new BasicTsKvEntry(ts, new LongDataEntry(key, 23L)); - assertThat(KvProtoUtil.fromProto(KvProtoUtil.toProto(ts, kvEntry))).as("deserialized").isEqualTo(kvEntry); + private static Stream attributeKvEntryData() { + return kvEntryData().map(kvEntry -> new BaseAttributeKvEntry(TS, kvEntry)); + } - kvEntry = new BasicTsKvEntry(ts, new DoubleDataEntry(key, 23.0)); - assertThat(KvProtoUtil.fromProto(KvProtoUtil.toProto(ts, kvEntry))).as("deserialized").isEqualTo(kvEntry); + private static List createTsKvEntryList(boolean withAggregation) { + return kvEntryData().map(kvEntry -> { + if (withAggregation) { + return new AggTsKvEntry(TS, kvEntry, 0); + } else { + return new BasicTsKvEntry(TS, kvEntry); + } + }).collect(Collectors.toList()); + } - kvEntry = new BasicTsKvEntry(ts, new StringDataEntry(key, "stringValue")); - assertThat(KvProtoUtil.fromProto(KvProtoUtil.toProto(ts, kvEntry))).as("deserialized").isEqualTo(kvEntry); + @ParameterizedTest + @EnumSource(DataType.class) + void protoDataTypeSerialization(DataType dataType) { + assertThat(KvProtoUtil.fromProto(KvProtoUtil.toProto(dataType))).as(dataType.name()).isEqualTo(dataType); + } - kvEntry = new BasicTsKvEntry(ts, new JsonDataEntry(key, "jsonValue")); - assertThat(KvProtoUtil.fromProto(KvProtoUtil.toProto(ts, kvEntry))).as("deserialized").isEqualTo(kvEntry); + @ParameterizedTest + @MethodSource("kvEntryData") + void protoKeyValueProtoSerialization(KvEntry kvEntry) { + assertThat(KvProtoUtil.fromProto(KvProtoUtil.toProto(kvEntry))) + .as("deserialized") + .isEqualTo(kvEntry); } - @Test - void protoListTsKvEntrySerialization() { - String key = "key"; - long ts = System.currentTimeMillis(); - KvEntry booleanDataEntry = new BooleanDataEntry(key, true); - KvEntry longDataEntry = new LongDataEntry(key, 23L); - KvEntry doubleDataEntry = new DoubleDataEntry(key, 23.0); - KvEntry stringDataEntry = new StringDataEntry(key, "stringValue"); - KvEntry jsonDataEntry = new JsonDataEntry(key, "jsonValue"); - List protoList = List.of( - new BasicTsKvEntry(ts, booleanDataEntry), - new BasicTsKvEntry(ts, longDataEntry), - new BasicTsKvEntry(ts, doubleDataEntry), - new BasicTsKvEntry(ts, stringDataEntry), - new BasicTsKvEntry(ts, jsonDataEntry) - ); - assertThat(KvProtoUtil.fromProtoList(KvProtoUtil.toProtoList(protoList))).as("deserialized").isEqualTo(protoList); + @ParameterizedTest + @MethodSource("basicTsKvEntryData") + void protoTsKvEntrySerialization(KvEntry kvEntry) { + assertThat(KvProtoUtil.fromProto(KvProtoUtil.toProto(TS, kvEntry))) + .as("deserialized") + .isEqualTo(kvEntry); + } - protoList = List.of( - new AggTsKvEntry(ts, booleanDataEntry, 3), - new AggTsKvEntry(ts, longDataEntry, 5), - new AggTsKvEntry(ts, doubleDataEntry, 2), - new AggTsKvEntry(ts, stringDataEntry, 1), - new AggTsKvEntry(ts, jsonDataEntry, 0) - ); - assertThat(KvProtoUtil.fromProtoList(KvProtoUtil.toProtoList(protoList))).as("deserialized").isEqualTo(protoList); + @ParameterizedTest + @ValueSource(booleans = {true, false}) + void protoListTsKvEntrySerialization(boolean withAggregation) { + List tsKvEntries = createTsKvEntryList(withAggregation); + assertThat(KvProtoUtil.fromProtoList(KvProtoUtil.toProtoList(tsKvEntries))) + .as("deserialized") + .isEqualTo(tsKvEntries); } @Test void protoListAttributeKvSerialization() { - String key = "key"; - long ts = System.currentTimeMillis(); - KvEntry booleanDataEntry = new BooleanDataEntry(key, true); - KvEntry longDataEntry = new LongDataEntry(key, 23L); - KvEntry doubleDataEntry = new DoubleDataEntry(key, 23.0); - KvEntry stringDataEntry = new StringDataEntry(key, "stringValue"); - KvEntry jsonDataEntry = new JsonDataEntry(key, "jsonValue"); - List protoList = List.of( - new BaseAttributeKvEntry(ts, booleanDataEntry), - new BaseAttributeKvEntry(ts, longDataEntry), - new BaseAttributeKvEntry(ts, doubleDataEntry), - new BaseAttributeKvEntry(ts, stringDataEntry), - new BaseAttributeKvEntry(ts, jsonDataEntry) - ); - assertThat(KvProtoUtil.toAttributeKvList(KvProtoUtil.attrToTsKvProtos(protoList))).as("deserialized").isEqualTo(protoList); + List protoList = attributeKvEntryData().toList(); + assertThat(KvProtoUtil.toAttributeKvList(KvProtoUtil.attrToTsKvProtos(protoList))) + .as("deserialized") + .isEqualTo(protoList); } } From 2e91612534cb6ccf46f311d0b5e9dc502e103bb2 Mon Sep 17 00:00:00 2001 From: Andrii Landiak Date: Mon, 25 Mar 2024 14:04:08 +0200 Subject: [PATCH 04/21] Minor refactoring --- .../subscription/TbSubscriptionUtils.java | 56 +------------------ .../server/common/util/KvProtoUtil.java | 34 +++++++++++ .../server/common/util/KvProtoUtilTest.java | 34 ++++++----- 3 files changed, 56 insertions(+), 68 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbSubscriptionUtils.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbSubscriptionUtils.java index 243a840f51..2d8440fa20 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/TbSubscriptionUtils.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbSubscriptionUtils.java @@ -23,13 +23,6 @@ import org.thingsboard.server.common.data.id.EntityIdFactory; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.UserId; import org.thingsboard.server.common.data.kv.AttributeKvEntry; -import org.thingsboard.server.common.data.kv.BasicTsKvEntry; -import org.thingsboard.server.common.data.kv.BooleanDataEntry; -import org.thingsboard.server.common.data.kv.DoubleDataEntry; -import org.thingsboard.server.common.data.kv.JsonDataEntry; -import org.thingsboard.server.common.data.kv.KvEntry; -import org.thingsboard.server.common.data.kv.LongDataEntry; -import org.thingsboard.server.common.data.kv.StringDataEntry; import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; import org.thingsboard.server.gen.transport.TransportProtos; @@ -55,8 +48,9 @@ import java.util.Map; import java.util.TreeMap; import java.util.UUID; -import static org.thingsboard.server.common.util.KvProtoUtil.fromKeyValueTypeProto; -import static org.thingsboard.server.common.util.KvProtoUtil.toKeyValueTypeProto; +import static org.thingsboard.server.common.util.KvProtoUtil.fromTsValueProtoList; +import static org.thingsboard.server.common.util.KvProtoUtil.toTsKvProtoBuilder; +import static org.thingsboard.server.common.util.KvProtoUtil.toTsValueProto; public class TbSubscriptionUtils { @@ -332,48 +326,4 @@ public class TbSubscriptionUtils { return ToCoreNotificationMsg.newBuilder().setToLocalSubscriptionServiceMsg(result).build(); } - public static TransportProtos.TsKvProto.Builder toTsKvProtoBuilder(long ts, KvEntry attr) { - TransportProtos.KeyValueProto.Builder dataBuilder = TransportProtos.KeyValueProto.newBuilder(); - dataBuilder.setKey(attr.getKey()); - dataBuilder.setType(toKeyValueTypeProto(attr.getDataType())); - switch (attr.getDataType()) { - case BOOLEAN -> attr.getBooleanValue().ifPresent(dataBuilder::setBoolV); - case LONG -> attr.getLongValue().ifPresent(dataBuilder::setLongV); - case DOUBLE -> attr.getDoubleValue().ifPresent(dataBuilder::setDoubleV); - case JSON -> attr.getJsonValue().ifPresent(dataBuilder::setJsonV); - case STRING -> attr.getStrValue().ifPresent(dataBuilder::setStringV); - } - return TransportProtos.TsKvProto.newBuilder().setTs(ts).setKv(dataBuilder); - } - - public static TransportProtos.TsValueProto toTsValueProto(long ts, KvEntry attr) { - TransportProtos.TsValueProto.Builder dataBuilder = TransportProtos.TsValueProto.newBuilder(); - dataBuilder.setTs(ts); - dataBuilder.setType(toKeyValueTypeProto(attr.getDataType())); - switch (attr.getDataType()) { - case BOOLEAN -> attr.getBooleanValue().ifPresent(dataBuilder::setBoolV); - case LONG -> attr.getLongValue().ifPresent(dataBuilder::setLongV); - case DOUBLE -> attr.getDoubleValue().ifPresent(dataBuilder::setDoubleV); - case JSON -> attr.getJsonValue().ifPresent(dataBuilder::setJsonV); - case STRING -> attr.getStrValue().ifPresent(dataBuilder::setStringV); - } - return dataBuilder.build(); - } - - private static List fromTsValueProtoList(String key, List dataList) { - List result = new ArrayList<>(dataList.size()); - dataList.forEach(proto -> result.add(new BasicTsKvEntry(proto.getTs(), fromTsKvProto(key, proto)))); - return result; - } - - private static KvEntry fromTsKvProto(String key, TransportProtos.TsValueProto proto) { - return switch (fromKeyValueTypeProto(proto.getType())) { - case BOOLEAN -> new BooleanDataEntry(key, proto.getBoolV()); - case LONG -> new LongDataEntry(key, proto.getLongV()); - case DOUBLE -> new DoubleDataEntry(key, proto.getDoubleV()); - case STRING -> new StringDataEntry(key, proto.getStringV()); - case JSON -> new JsonDataEntry(key, proto.getJsonV()); - }; - } - } diff --git a/common/proto/src/main/java/org/thingsboard/server/common/util/KvProtoUtil.java b/common/proto/src/main/java/org/thingsboard/server/common/util/KvProtoUtil.java index 14368693c5..74674e1e45 100644 --- a/common/proto/src/main/java/org/thingsboard/server/common/util/KvProtoUtil.java +++ b/common/proto/src/main/java/org/thingsboard/server/common/util/KvProtoUtil.java @@ -114,6 +114,40 @@ public class KvProtoUtil { }; } + public static TransportProtos.TsKvProto.Builder toTsKvProtoBuilder(long ts, KvEntry kvEntry) { + return TransportProtos.TsKvProto.newBuilder().setTs(ts).setKv(KvProtoUtil.toKeyValueTypeProto(kvEntry)); + } + + public static List fromTsValueProtoList(String key, List dataList) { + List result = new ArrayList<>(dataList.size()); + dataList.forEach(proto -> result.add(new BasicTsKvEntry(proto.getTs(), fromTsValueProto(key, proto)))); + return result; + } + + public static TransportProtos.TsValueProto toTsValueProto(long ts, KvEntry attr) { + TransportProtos.TsValueProto.Builder dataBuilder = TransportProtos.TsValueProto.newBuilder(); + dataBuilder.setTs(ts); + dataBuilder.setType(toKeyValueTypeProto(attr.getDataType())); + switch (attr.getDataType()) { + case BOOLEAN -> attr.getBooleanValue().ifPresent(dataBuilder::setBoolV); + case LONG -> attr.getLongValue().ifPresent(dataBuilder::setLongV); + case DOUBLE -> attr.getDoubleValue().ifPresent(dataBuilder::setDoubleV); + case JSON -> attr.getJsonValue().ifPresent(dataBuilder::setJsonV); + case STRING -> attr.getStrValue().ifPresent(dataBuilder::setStringV); + } + return dataBuilder.build(); + } + + public static KvEntry fromTsValueProto(String key, TransportProtos.TsValueProto proto) { + return switch (fromKeyValueTypeProto(proto.getType())) { + case BOOLEAN -> new BooleanDataEntry(key, proto.getBoolV()); + case LONG -> new LongDataEntry(key, proto.getLongV()); + case DOUBLE -> new DoubleDataEntry(key, proto.getDoubleV()); + case STRING -> new StringDataEntry(key, proto.getStringV()); + case JSON -> new JsonDataEntry(key, proto.getJsonV()); + }; + } + public static TransportProtos.KeyValueType toKeyValueTypeProto(DataType dataType) { return TransportProtos.KeyValueType.forNumber(dataType.getProtoNumber()); } diff --git a/common/proto/src/test/java/org/thingsboard/server/common/util/KvProtoUtilTest.java b/common/proto/src/test/java/org/thingsboard/server/common/util/KvProtoUtilTest.java index 1fbfaf52d9..72d71c7c93 100644 --- a/common/proto/src/test/java/org/thingsboard/server/common/util/KvProtoUtilTest.java +++ b/common/proto/src/test/java/org/thingsboard/server/common/util/KvProtoUtilTest.java @@ -15,7 +15,6 @@ */ package org.thingsboard.server.common.util; -import org.junit.jupiter.api.Test; import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.provider.EnumSource; import org.junit.jupiter.params.provider.MethodSource; @@ -58,8 +57,8 @@ class KvProtoUtilTest { return kvEntryData().map(kvEntry -> new BasicTsKvEntry(TS, kvEntry)); } - private static Stream attributeKvEntryData() { - return kvEntryData().map(kvEntry -> new BaseAttributeKvEntry(TS, kvEntry)); + private static Stream> attributeKvEntryData() { + return Stream.of(kvEntryData().map(kvEntry -> new BaseAttributeKvEntry(TS, kvEntry)).toList()); } private static List createTsKvEntryList(boolean withAggregation) { @@ -75,23 +74,29 @@ class KvProtoUtilTest { @ParameterizedTest @EnumSource(DataType.class) void protoDataTypeSerialization(DataType dataType) { - assertThat(KvProtoUtil.fromKeyValueTypeProto(KvProtoUtil.toKeyValueTypeProto(dataType))).as(dataType.name()).isEqualTo(dataType); + assertThat(KvProtoUtil.fromKeyValueTypeProto(KvProtoUtil.toKeyValueTypeProto(dataType))) + .as(dataType.name()).isEqualTo(dataType); } @ParameterizedTest @MethodSource("kvEntryData") void protoKeyValueProtoSerialization(KvEntry kvEntry) { assertThat(KvProtoUtil.fromTsKvProto(KvProtoUtil.toKeyValueTypeProto(kvEntry))) - .as("deserialized") - .isEqualTo(kvEntry); + .as("deserialized").isEqualTo(kvEntry); } @ParameterizedTest @MethodSource("basicTsKvEntryData") void protoTsKvEntrySerialization(KvEntry kvEntry) { assertThat(KvProtoUtil.fromTsKvProto(KvProtoUtil.toTsKvProto(TS, kvEntry))) - .as("deserialized") - .isEqualTo(kvEntry); + .as("deserialized").isEqualTo(kvEntry); + } + + @ParameterizedTest + @MethodSource("kvEntryData") + void protoTsValueSerialization(KvEntry kvEntry) { + assertThat(KvProtoUtil.fromTsValueProto(kvEntry.getKey(), KvProtoUtil.toTsValueProto(TS, kvEntry))) + .as("deserialized").isEqualTo(kvEntry); } @ParameterizedTest @@ -99,16 +104,15 @@ class KvProtoUtilTest { void protoListTsKvEntrySerialization(boolean withAggregation) { List tsKvEntries = createTsKvEntryList(withAggregation); assertThat(KvProtoUtil.fromTsKvProtoList(KvProtoUtil.toTsKvProtoList(tsKvEntries))) - .as("deserialized") - .isEqualTo(tsKvEntries); + .as("deserialized").isEqualTo(tsKvEntries); } - @Test - void protoListAttributeKvSerialization() { - List protoList = attributeKvEntryData().toList(); - assertThat(KvProtoUtil.toAttributeKvList(KvProtoUtil.attrToTsKvProtos(protoList))) + @ParameterizedTest + @MethodSource("attributeKvEntryData") + void protoListAttributeKvSerialization(List attributeKvEntries) { + assertThat(KvProtoUtil.toAttributeKvList(KvProtoUtil.attrToTsKvProtos(attributeKvEntries))) .as("deserialized") - .isEqualTo(protoList); + .isEqualTo(attributeKvEntries); } } From eba645c542e4f06ff67ce8f3bf7154f3ba4e51fc Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Tue, 26 Mar 2024 19:25:44 +0100 Subject: [PATCH 05/21] Nashorn LOCAL_JS_SANDBOX_MAX_MEMORY introduced --- application/src/main/resources/thingsboard.yml | 2 ++ .../script/api/AbstractScriptInvokeService.java | 4 ++-- .../thingsboard/script/api/js/NashornJsInvokeService.java | 7 +++++-- 3 files changed, 9 insertions(+), 4 deletions(-) diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index daf9f5d1c1..01f0d83771 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -866,6 +866,8 @@ js: monitor_thread_pool_size: "${LOCAL_JS_SANDBOX_MONITOR_THREAD_POOL_SIZE:4}" # Maximum CPU time in milliseconds allowed for script execution max_cpu_time: "${LOCAL_JS_SANDBOX_MAX_CPU_TIME:8000}" + # Maximum memory in Bytes which JS executor thread can allocate (approximate calculation). A zero memory limit in combination with a non-zero CPU limit is not recommended due to the implementation of Nashorn 0.4.2. 100MiB is effectively unlimited for most cases + max_memory: "${LOCAL_JS_SANDBOX_MAX_MEMORY:104857600}" # Maximum allowed JavaScript execution errors before JavaScript will be blacklisted max_errors: "${LOCAL_JS_SANDBOX_MAX_ERRORS:3}" # JS Eval max request timeout. 0 - no timeout diff --git a/common/script/script-api/src/main/java/org/thingsboard/script/api/AbstractScriptInvokeService.java b/common/script/script-api/src/main/java/org/thingsboard/script/api/AbstractScriptInvokeService.java index f88625f7fd..32e209b2ea 100644 --- a/common/script/script-api/src/main/java/org/thingsboard/script/api/AbstractScriptInvokeService.java +++ b/common/script/script-api/src/main/java/org/thingsboard/script/api/AbstractScriptInvokeService.java @@ -145,14 +145,14 @@ public abstract class AbstractScriptInvokeService implements ScriptInvokeService log.trace("[{}] InvokeScript uuid {} with timeout {}ms", tenantId, scriptId, getMaxInvokeRequestsTimeout()); var task = doInvokeFunction(scriptId, args); - var resultFuture = Futures.transformAsync(task.getResultFuture(), output -> { + var resultFuture = Futures.transform(task.getResultFuture(), output -> { String result = JacksonUtil.toString(output); if (resultSizeExceeded(result)) { throw new TbScriptException(scriptId, TbScriptException.ErrorCode.OTHER, null, new RuntimeException( format("Script invocation result exceeds maximum allowed size of %s symbols", getMaxResultSize()) )); } - return Futures.immediateFuture(output); + return output; }, MoreExecutors.directExecutor()); return withTimeoutAndStatsCallback(scriptId, task, resultFuture, invokeCallback, getMaxInvokeRequestsTimeout()); diff --git a/common/script/script-api/src/main/java/org/thingsboard/script/api/js/NashornJsInvokeService.java b/common/script/script-api/src/main/java/org/thingsboard/script/api/js/NashornJsInvokeService.java index 0e37bd89d6..6a07e596ef 100644 --- a/common/script/script-api/src/main/java/org/thingsboard/script/api/js/NashornJsInvokeService.java +++ b/common/script/script-api/src/main/java/org/thingsboard/script/api/js/NashornJsInvokeService.java @@ -41,7 +41,6 @@ import java.util.Optional; import java.util.UUID; import java.util.concurrent.Executor; import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; import java.util.concurrent.locks.ReentrantLock; @Slf4j @@ -65,6 +64,9 @@ public class NashornJsInvokeService extends AbstractJsInvokeService { @Value("${js.local.max_cpu_time}") private long maxCpuTime; + @Value("${js.local.max_memory}") + private long maxMemory; + @Getter @Value("${js.local.max_errors}") private int maxErrors; @@ -107,12 +109,13 @@ public class NashornJsInvokeService extends AbstractJsInvokeService { @Override public void init() { super.init(); - jsExecutor = MoreExecutors.listeningDecorator(Executors.newWorkStealingPool(jsExecutorThreadPoolSize)); + jsExecutor = MoreExecutors.listeningDecorator(ThingsBoardExecutors.newWorkStealingPool(jsExecutorThreadPoolSize, "nashorn-js-executor")); if (useJsSandbox) { sandbox = NashornSandboxes.create(); monitorExecutorService = ThingsBoardExecutors.newWorkStealingPool(monitorThreadPoolSize, "nashorn-js-monitor"); sandbox.setExecutor(monitorExecutorService); sandbox.setMaxCPUTime(maxCpuTime); + sandbox.setMaxMemory(maxMemory); sandbox.allowNoBraces(false); sandbox.allowLoadFunctions(true); sandbox.setMaxPreparedStatements(30); From 3d1f1e4af0ccc844fd9c752af1db26df99a943ee Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Tue, 26 Mar 2024 19:27:23 +0100 Subject: [PATCH 06/21] NashornJsInvokeServiceTest refactored --- .../script/NashornJsInvokeServiceTest.java | 31 +++++++++++-------- .../src/test/resources/logback-test.xml | 2 +- 2 files changed, 19 insertions(+), 14 deletions(-) diff --git a/application/src/test/java/org/thingsboard/server/service/script/NashornJsInvokeServiceTest.java b/application/src/test/java/org/thingsboard/server/service/script/NashornJsInvokeServiceTest.java index 55c521c309..bb9dc57d3c 100644 --- a/application/src/test/java/org/thingsboard/server/service/script/NashornJsInvokeServiceTest.java +++ b/application/src/test/java/org/thingsboard/server/service/script/NashornJsInvokeServiceTest.java @@ -15,13 +15,13 @@ */ package org.thingsboard.server.service.script; -import com.fasterxml.jackson.databind.node.ObjectNode; +import lombok.extern.slf4j.Slf4j; import org.junit.Assert; import org.junit.jupiter.api.Test; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; import org.springframework.test.context.TestPropertySource; -import org.thingsboard.common.util.JacksonUtil; +import org.thingsboard.common.util.TbStopWatch; import org.thingsboard.script.api.ScriptType; import org.thingsboard.script.api.js.NashornJsInvokeService; import org.thingsboard.server.common.data.id.TenantId; @@ -32,6 +32,7 @@ import java.util.UUID; import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; +import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.thingsboard.server.common.data.msg.TbMsgType.POST_TELEMETRY_REQUEST; @@ -41,8 +42,9 @@ import static org.thingsboard.server.common.data.msg.TbMsgType.POST_TELEMETRY_RE "js.max_script_body_size=50", "js.max_total_args_size=50", "js.max_result_size=50", - "js.local.max_errors=2" + "js.local.max_errors=2", }) +@Slf4j class NashornJsInvokeServiceTest extends AbstractControllerTest { @Autowired @@ -56,23 +58,26 @@ class NashornJsInvokeServiceTest extends AbstractControllerTest { int iterations = 1000; UUID scriptId = evalScript("return msg.temperature > 20"); // warmup - ObjectNode msg = JacksonUtil.newObjectNode(); - for (int i = 0; i < 100; i++) { - msg.put("temperature", i); + log.info("Warming up 1000 times..."); + var warmupWatch = TbStopWatch.create(); + for (int i = 0; i < 1000; i++) { boolean expected = i > 20; - boolean result = Boolean.valueOf(invokeScript(scriptId, JacksonUtil.toString(msg))); + boolean result = Boolean.parseBoolean(invokeScript(scriptId, "{\"temperature\":" + i + "}")); Assert.assertEquals(expected, result); } - long startTs = System.currentTimeMillis(); + log.info("Warming up finished in {} ms", warmupWatch.stopAndGetTotalTimeMillis()); + log.info("Starting performance test..."); + var watch = TbStopWatch.create(); for (int i = 0; i < iterations; i++) { - msg.put("temperature", i); boolean expected = i > 20; - boolean result = Boolean.valueOf(invokeScript(scriptId, JacksonUtil.toString(msg))); + boolean result = Boolean.parseBoolean(invokeScript(scriptId, "{\"temperature\":" + i + "}")); + log.debug("asserting result"); Assert.assertEquals(expected, result); } - long duration = System.currentTimeMillis() - startTs; - System.out.println(iterations + " invocations took: " + duration + "ms"); - Assert.assertTrue(duration < TimeUnit.MINUTES.toMillis(4)); + long duration = watch.stopAndGetTotalTimeMillis(); + log.info("Performance test with {} invocations took: {} ms", iterations, duration); + assertThat(duration).as("duration ms") + .isLessThan(TimeUnit.MINUTES.toMillis(1)); // effective exec time is about 500ms } @Test diff --git a/application/src/test/resources/logback-test.xml b/application/src/test/resources/logback-test.xml index 981bcab132..d72bccb7a6 100644 --- a/application/src/test/resources/logback-test.xml +++ b/application/src/test/resources/logback-test.xml @@ -16,7 +16,7 @@ - + From cd722a11484bf3f54cdeb628135af7018ac36f65 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Tue, 26 Mar 2024 19:34:13 +0100 Subject: [PATCH 07/21] bump delight-nashorn-sandbox.version to the latest 0.4.2. As it supports Java 17 --- pom.xml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pom.xml b/pom.xml index d2450e0a65..f8e6c1ec1f 100755 --- a/pom.xml +++ b/pom.xml @@ -103,7 +103,7 @@ org/thingsboard/server/extensions/core/plugin/telemetry/gen/**/* 5.0.2 - 0.2.1 + 0.4.2 15.4