From 8af37beb4a7df2d51cbca5a8654329939f2db34f Mon Sep 17 00:00:00 2001 From: dshvaika Date: Tue, 20 May 2025 16:19:44 +0300 Subject: [PATCH 01/10] New Cassandra rate limits: separated for Read and Write + Core and Rule Engine --- .../service/limits/RateLimitServiceTest.java | 12 +- .../server/common/data/limit/LimitedApi.java | 13 +- .../common/data/limit/LimitedApiEntry.java | 34 ++++++ .../common/data/limit/LimitedApiUtil.java | 67 +++++++++++ .../DefaultTenantProfileConfiguration.java | 6 +- .../common/data/limit/LimitedApiUtilTest.java | 113 ++++++++++++++++++ .../server/common/msg/tools/TbRateLimits.java | 27 +++-- .../common/msg/tools/TbRateLimitsTest.java | 63 ++++++++++ .../DefaultTbServiceInfoProvider.java | 13 +- .../discovery/TbServiceInfoProvider.java | 2 + dao/pom.xml | 4 + .../CassandraBufferedRateReadExecutor.java | 13 +- .../CassandraBufferedRateWriteExecutor.java | 13 +- .../util/AbstractBufferedRateExecutor.java | 24 +++- .../dao/util/BufferedRateExecutorType.java | 40 +++++++ 15 files changed, 407 insertions(+), 37 deletions(-) create mode 100644 common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApiEntry.java create mode 100644 common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApiUtil.java create mode 100644 common/data/src/test/java/org/thingsboard/server/common/data/limit/LimitedApiUtilTest.java create mode 100644 common/message/src/test/java/org/thingsboard/server/common/msg/tools/TbRateLimitsTest.java create mode 100644 dao/src/main/java/org/thingsboard/server/dao/util/BufferedRateExecutorType.java diff --git a/application/src/test/java/org/thingsboard/server/service/limits/RateLimitServiceTest.java b/application/src/test/java/org/thingsboard/server/service/limits/RateLimitServiceTest.java index 3eb5793c74..f85ad848ab 100644 --- a/application/src/test/java/org/thingsboard/server/service/limits/RateLimitServiceTest.java +++ b/application/src/test/java/org/thingsboard/server/service/limits/RateLimitServiceTest.java @@ -52,7 +52,7 @@ public class RateLimitServiceTest { public void beforeEach() { tenantProfileCache = Mockito.mock(DefaultTbTenantProfileCache.class); rateLimitService = new DefaultRateLimitService(tenantProfileCache, mock(NotificationRuleProcessor.class), 60, 100); - tenantId = new TenantId(UUID.randomUUID()); + tenantId = TenantId.fromUUID(UUID.randomUUID()); } @Test @@ -67,7 +67,10 @@ public class RateLimitServiceTest { profileConfiguration.setTenantServerRestLimitsConfiguration(rateLimit); profileConfiguration.setCustomerServerRestLimitsConfiguration(rateLimit); profileConfiguration.setWsUpdatesPerSessionRateLimit(rateLimit); - profileConfiguration.setCassandraQueryTenantRateLimitsConfiguration(rateLimit); + profileConfiguration.setCassandraReadQueryTenantCoreRateLimits(rateLimit); + profileConfiguration.setCassandraWriteQueryTenantCoreRateLimits(rateLimit); + profileConfiguration.setCassandraReadQueryTenantRuleEngineRateLimits(rateLimit); + profileConfiguration.setCassandraWriteQueryTenantRuleEngineRateLimits(rateLimit); profileConfiguration.setEdgeEventRateLimits(rateLimit); profileConfiguration.setEdgeEventRateLimitsPerEdge(rateLimit); profileConfiguration.setEdgeUplinkMessagesRateLimits(rateLimit); @@ -79,7 +82,10 @@ public class RateLimitServiceTest { LimitedApi.ENTITY_IMPORT, LimitedApi.NOTIFICATION_REQUESTS, LimitedApi.REST_REQUESTS_PER_CUSTOMER, - LimitedApi.CASSANDRA_QUERIES, + LimitedApi.CASSANDRA_READ_QUERIES_CORE, + LimitedApi.CASSANDRA_WRITE_QUERIES_CORE, + LimitedApi.CASSANDRA_READ_QUERIES_RULE_ENGINE, + LimitedApi.CASSANDRA_WRITE_QUERIES_RULE_ENGINE, LimitedApi.EDGE_EVENTS, LimitedApi.EDGE_EVENTS_PER_EDGE, LimitedApi.EDGE_UPLINK_MESSAGES, diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApi.java b/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApi.java index db7f14171b..766a1f4eb3 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApi.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApi.java @@ -30,7 +30,18 @@ public enum LimitedApi { REST_REQUESTS_PER_TENANT(DefaultTenantProfileConfiguration::getTenantServerRestLimitsConfiguration, "REST API requests", true), REST_REQUESTS_PER_CUSTOMER(DefaultTenantProfileConfiguration::getCustomerServerRestLimitsConfiguration, "REST API requests per customer", false), WS_UPDATES_PER_SESSION(DefaultTenantProfileConfiguration::getWsUpdatesPerSessionRateLimit, "WS updates per session", true), - CASSANDRA_QUERIES(DefaultTenantProfileConfiguration::getCassandraQueryTenantRateLimitsConfiguration, "Cassandra queries", true), + CASSANDRA_WRITE_QUERIES_CORE(DefaultTenantProfileConfiguration::getCassandraReadQueryTenantCoreRateLimits, "Rest API and WS telemetry read queries", true), + CASSANDRA_READ_QUERIES_CORE(DefaultTenantProfileConfiguration::getCassandraWriteQueryTenantCoreRateLimits, "Rest API and WS telemetry write queries", true), + CASSANDRA_WRITE_QUERIES_RULE_ENGINE(DefaultTenantProfileConfiguration::getCassandraReadQueryTenantRuleEngineRateLimits, "Rule Engine telemetry read queries", true), + CASSANDRA_READ_QUERIES_RULE_ENGINE(DefaultTenantProfileConfiguration::getCassandraWriteQueryTenantRuleEngineRateLimits, "Rule Engine telemetry write queries", true), + CASSANDRA_READ_QUERIES_MONOLITH( + LimitedApiUtil.merge( + DefaultTenantProfileConfiguration::getCassandraReadQueryTenantCoreRateLimits, + DefaultTenantProfileConfiguration::getCassandraReadQueryTenantRuleEngineRateLimits), "Telemetry read queries", true), + CASSANDRA_WRITE_QUERIES_MONOLITH( + LimitedApiUtil.merge( + DefaultTenantProfileConfiguration::getCassandraWriteQueryTenantCoreRateLimits, + DefaultTenantProfileConfiguration::getCassandraWriteQueryTenantRuleEngineRateLimits), "Telemetry write queries", true), EDGE_EVENTS(DefaultTenantProfileConfiguration::getEdgeEventRateLimits, "Edge events", true), EDGE_EVENTS_PER_EDGE(DefaultTenantProfileConfiguration::getEdgeEventRateLimitsPerEdge, "Edge events per edge", false), EDGE_UPLINK_MESSAGES(DefaultTenantProfileConfiguration::getEdgeUplinkMessagesRateLimits, "Edge uplink messages", true), diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApiEntry.java b/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApiEntry.java new file mode 100644 index 0000000000..3082b521b6 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApiEntry.java @@ -0,0 +1,34 @@ +/** + * Copyright © 2016-2025 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.data.limit; + +public record LimitedApiEntry(long capacity, long durationSeconds) { + + public static LimitedApiEntry parse(String s) { + String[] parts = s.split(":"); + return new LimitedApiEntry(Long.parseLong(parts[0]), Long.parseLong(parts[1])); + } + + public double rps() { + return (double) capacity / durationSeconds; + } + + @Override + public String toString() { + return capacity + ":" + durationSeconds; + } + +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApiUtil.java b/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApiUtil.java new file mode 100644 index 0000000000..a90c7eb825 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApiUtil.java @@ -0,0 +1,67 @@ +/** + * Copyright © 2016-2025 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.data.limit; + +import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration; + +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.function.Function; +import java.util.stream.Collectors; + +public class LimitedApiUtil { + + public static List parseConfig(String config) { + if (config == null || config.isEmpty()) { + return Collections.emptyList(); + } + return Arrays.stream(config.split(",")) + .map(LimitedApiEntry::parse) + .toList(); + } + + public static Function merge( + Function configExtractor1, + Function configExtractor2) { + return config -> { + String config1 = configExtractor1.apply(config); + String config2 = configExtractor2.apply(config); + return LimitedApiUtil.mergeStrConfigs(config1, config2); // merges the configs + }; + } + + private static String mergeStrConfigs(String firstConfig, String secondConfig) { + List all = new ArrayList<>(); + all.addAll(parseConfig(firstConfig)); + all.addAll(parseConfig(secondConfig)); + + Map merged = new HashMap<>(); + + for (LimitedApiEntry entry : all) { + merged.merge(entry.durationSeconds(), entry.capacity(), Long::sum); + } + + return merged.entrySet().stream() + .sorted(Map.Entry.comparingByKey()) // optional: sort by duration + .map(e -> e.getValue() + ":" + e.getKey()) + .collect(Collectors.joining(",")); + } + +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java index f256b02d9a..d61f48b148 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java @@ -121,7 +121,11 @@ public class DefaultTenantProfileConfiguration implements TenantProfileConfigura private long maxWsSubscriptionsPerPublicUser; private String wsUpdatesPerSessionRateLimit; - private String cassandraQueryTenantRateLimitsConfiguration; + private String cassandraReadQueryTenantCoreRateLimits; + private String cassandraWriteQueryTenantCoreRateLimits; + + private String cassandraReadQueryTenantRuleEngineRateLimits; + private String cassandraWriteQueryTenantRuleEngineRateLimits; private String edgeEventRateLimits; private String edgeEventRateLimitsPerEdge; diff --git a/common/data/src/test/java/org/thingsboard/server/common/data/limit/LimitedApiUtilTest.java b/common/data/src/test/java/org/thingsboard/server/common/data/limit/LimitedApiUtilTest.java new file mode 100644 index 0000000000..6486de3221 --- /dev/null +++ b/common/data/src/test/java/org/thingsboard/server/common/data/limit/LimitedApiUtilTest.java @@ -0,0 +1,113 @@ +/** + * Copyright © 2016-2025 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.data.limit; + +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; +import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration; + +import java.util.List; +import java.util.function.Function; + +import static org.assertj.core.api.Assertions.assertThat; + +class LimitedApiUtilTest { + + @Test + @DisplayName("LimitedApiUtil should parse single entry correctly") + void testParseSingleEntry() { + List entries = LimitedApiUtil.parseConfig("100:60"); + + assertThat(entries).hasSize(1); + assertThat(entries.get(0).capacity()).isEqualTo(100); + assertThat(entries.get(0).durationSeconds()).isEqualTo(60); + } + + @Test + @DisplayName("LimitedApiUtil should parse multiple entries correctly") + void testParseMultipleEntries() { + List entries = LimitedApiUtil.parseConfig("100:60,200:30"); + + assertThat(entries).hasSize(2); + assertThat(entries.get(0).capacity()).isEqualTo(100); + assertThat(entries.get(0).durationSeconds()).isEqualTo(60); + assertThat(entries.get(1).capacity()).isEqualTo(200); + assertThat(entries.get(1).durationSeconds()).isEqualTo(30); + } + + @Test + @DisplayName("LimitedApiUtil should return empty list for null or empty config") + void testParseEmptyConfig() { + assertThat(LimitedApiUtil.parseConfig(null)).isEmpty(); + assertThat(LimitedApiUtil.parseConfig("")).isEmpty(); + } + + @Test + @DisplayName("LimitedApiUtil should merge two configs by summing capacities with same durations") + void testMergeStrConfigs() { + Function extractor1 = cfg -> "100:60,50:30"; + Function extractor2 = cfg -> "200:60,25:10"; + + // Fake config instance (not used directly in lambda logic) + DefaultTenantProfileConfiguration config = new DefaultTenantProfileConfiguration(); + + String result = LimitedApiUtil.merge(extractor1, extractor2).apply(config); + + // Should be: 300:60 (100+200), 50:30, 25:10 + assertThat(result).isEqualTo("25:10,50:30,300:60"); + } + + @Test + @DisplayName("LimitedApiUtil should merge configs when one is empty") + void testMergeWithEmptyOne() { + Function extractor1 = cfg -> "100:60"; + Function extractor2 = cfg -> ""; + + // Fake config instance (not used directly in lambda logic) + DefaultTenantProfileConfiguration config = new DefaultTenantProfileConfiguration(); + String result = LimitedApiUtil.merge(extractor1, extractor2).apply(config); + + assertThat(result).isEqualTo("100:60"); + } + + @Test + @DisplayName("LimitedApiUtil should merge configs when both have distinct durations") + void testMergeWithDistinctDurations() { + Function extractor1 = cfg -> "100:60"; + Function extractor2 = cfg -> "200:10"; + + // Fake config instance (not used directly in lambda logic) + DefaultTenantProfileConfiguration config = new DefaultTenantProfileConfiguration(); + String result = LimitedApiUtil.merge(extractor1, extractor2).apply(config); + + assertThat(result).isEqualTo("200:10,100:60"); + } + + @Test + @DisplayName("LimitedApiUtil shouldn't have duplicate durations in the same config!") + void testMergeHandlesDuplicatesInSingleConfig() { + Function extractor1 = cfg -> "100:60,200:60"; + Function extractor2 = cfg -> ""; + + // Fake config instance (not used directly in lambda logic) + DefaultTenantProfileConfiguration config = new DefaultTenantProfileConfiguration(); + String result = LimitedApiUtil.merge(extractor1, extractor2).apply(config); + + // 100+200 = 300 for duration 60. Currently possible to save the same "per seconds" config from the UI. + // This must be fixed, so we will merge only two different rate limits. + assertThat(result).isEqualTo("300:60"); + } +} diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/tools/TbRateLimits.java b/common/message/src/main/java/org/thingsboard/server/common/msg/tools/TbRateLimits.java index 8e8933c324..1e2d8dd4ce 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/tools/TbRateLimits.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/tools/TbRateLimits.java @@ -16,13 +16,17 @@ package org.thingsboard.server.common.msg.tools; import io.github.bucket4j.Bandwidth; +import io.github.bucket4j.BandwidthBuilder; import io.github.bucket4j.Bucket; import io.github.bucket4j.Refill; import io.github.bucket4j.local.LocalBucket; import io.github.bucket4j.local.LocalBucketBuilder; import lombok.Getter; +import org.thingsboard.server.common.data.limit.LimitedApiEntry; +import org.thingsboard.server.common.data.limit.LimitedApiUtil; import java.time.Duration; +import java.util.List; /** * Created by ashvayka on 22.10.18. @@ -38,20 +42,19 @@ public class TbRateLimits { } public TbRateLimits(String limitsConfiguration, boolean refillIntervally) { - LocalBucketBuilder builder = Bucket.builder(); - boolean initialized = false; - for (String limitSrc : limitsConfiguration.split(",")) { - long capacity = Long.parseLong(limitSrc.split(":")[0]); - long duration = Long.parseLong(limitSrc.split(":")[1]); - Refill refill = refillIntervally ? Refill.intervally(capacity, Duration.ofSeconds(duration)) : Refill.greedy(capacity, Duration.ofSeconds(duration)); - builder.addLimit(Bandwidth.classic(capacity, refill)); - initialized = true; - } - if (initialized) { - bucket = builder.build(); - } else { + List limitedApiEntries = LimitedApiUtil.parseConfig(limitsConfiguration); + if (limitedApiEntries.isEmpty()) { throw new IllegalArgumentException("Failed to parse rate limits configuration: " + limitsConfiguration); } + LocalBucketBuilder localBucket = Bucket.builder(); + for (LimitedApiEntry entry : limitedApiEntries) { + BandwidthBuilder.BandwidthBuilderRefillStage bandwidthBuilder = Bandwidth.builder().capacity(entry.capacity()); + Bandwidth bandwidth = refillIntervally ? + bandwidthBuilder.refillIntervally(entry.capacity(), Duration.ofSeconds(entry.durationSeconds())).build() : + bandwidthBuilder.refillGreedy(entry.capacity(), Duration.ofSeconds(entry.durationSeconds())).build(); + localBucket.addLimit(bandwidth); + } + this.bucket = localBucket.build(); this.configuration = limitsConfiguration; } diff --git a/common/message/src/test/java/org/thingsboard/server/common/msg/tools/TbRateLimitsTest.java b/common/message/src/test/java/org/thingsboard/server/common/msg/tools/TbRateLimitsTest.java new file mode 100644 index 0000000000..a6a95da9c1 --- /dev/null +++ b/common/message/src/test/java/org/thingsboard/server/common/msg/tools/TbRateLimitsTest.java @@ -0,0 +1,63 @@ +/** + * Copyright © 2016-2025 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.msg.tools; + +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +class TbRateLimitsTest { + + @Test + @DisplayName("TbRateLimits should construct with single rate limit") + void testSingleLimitConstructor() { + TbRateLimits limits = new TbRateLimits("10:1", false); + assertThat(limits.getConfiguration()).isEqualTo("10:1"); + } + + @Test + @DisplayName("TbRateLimits should construct with multiple rate limits") + void testMultipleLimitConstructor() { + String config = "10:1,100:10"; + TbRateLimits limits = new TbRateLimits(config, false); + assertThat(limits.getConfiguration()).isEqualTo(config); + } + + @Test + @DisplayName("TbRateLimits should throw IllegalArgumentException on empty string") + void testEmptyConfigThrows() { + assertThatThrownBy(() -> new TbRateLimits("", false)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage("Failed to parse rate limits configuration: "); + } + + @Test + @DisplayName("TbRateLimits should throw NumberFormatException on malformed value") + void testMalformedConfigThrows() { + assertThatThrownBy(() -> new TbRateLimits("not_a_number:second", false)) + .isInstanceOf(NumberFormatException.class); + } + + @Test + @DisplayName("TbRateLimits should throw ArrayIndexOutOfBoundsException on missing colon") + void testColonMissingThrows() { + assertThatThrownBy(() -> new TbRateLimits("100", false)) + .isInstanceOf(ArrayIndexOutOfBoundsException.class); + } + +} \ No newline at end of file diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/DefaultTbServiceInfoProvider.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/DefaultTbServiceInfoProvider.java index 609d3f8eee..7ff6fafbca 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/DefaultTbServiceInfoProvider.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/DefaultTbServiceInfoProvider.java @@ -78,11 +78,9 @@ public class DefaultTbServiceInfoProvider implements TbServiceInfoProvider { } } log.info("Current Service ID: {}", serviceId); - if (serviceType.equalsIgnoreCase("monolith")) { - serviceTypes = List.of(ServiceType.values()); - } else { - serviceTypes = Collections.singletonList(ServiceType.of(serviceType)); - } + serviceTypes = isMonolith() ? + List.of(ServiceType.values()) : + Collections.singletonList(ServiceType.of(serviceType)); if (!serviceTypes.contains(ServiceType.TB_RULE_ENGINE) || assignedTenantProfiles == null) { assignedTenantProfiles = Collections.emptySet(); } @@ -113,6 +111,11 @@ public class DefaultTbServiceInfoProvider implements TbServiceInfoProvider { return serviceInfo; } + @Override + public boolean isMonolith() { + return serviceType.equalsIgnoreCase("monolith"); + } + @Override public boolean isService(ServiceType serviceType) { return serviceTypes.contains(serviceType); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/TbServiceInfoProvider.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/TbServiceInfoProvider.java index 51a6d808dd..a844440a21 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/TbServiceInfoProvider.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/TbServiceInfoProvider.java @@ -29,6 +29,8 @@ public interface TbServiceInfoProvider { ServiceInfo getServiceInfo(); + boolean isMonolith(); + boolean isService(ServiceType serviceType); ServiceInfo generateNewServiceInfoWithCurrentSystemInfo(); diff --git a/dao/pom.xml b/dao/pom.xml index 2206d8a7d7..67b5a36c92 100644 --- a/dao/pom.xml +++ b/dao/pom.xml @@ -59,6 +59,10 @@ org.thingsboard.common util + + org.thingsboard.common + queue + com.networknt json-schema-validator diff --git a/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraBufferedRateReadExecutor.java b/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraBufferedRateReadExecutor.java index 259cd908fc..a7aa0dc5dc 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraBufferedRateReadExecutor.java +++ b/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraBufferedRateReadExecutor.java @@ -28,7 +28,9 @@ import org.thingsboard.server.common.stats.StatsFactory; import org.thingsboard.server.dao.entity.EntityService; import org.thingsboard.server.dao.util.AbstractBufferedRateExecutor; import org.thingsboard.server.dao.util.AsyncTaskContext; +import org.thingsboard.server.dao.util.BufferedRateExecutorType; import org.thingsboard.server.dao.util.NoSqlAnyDao; +import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; /** * Created by ashvayka on 24.10.18. @@ -38,8 +40,6 @@ import org.thingsboard.server.dao.util.NoSqlAnyDao; @NoSqlAnyDao public class CassandraBufferedRateReadExecutor extends AbstractBufferedRateExecutor { - static final String BUFFER_NAME = "Read"; - public CassandraBufferedRateReadExecutor( @Value("${cassandra.query.buffer_size}") int queueLimit, @Value("${cassandra.query.concurrent_limit}") int concurrencyLimit, @@ -51,9 +51,10 @@ public class CassandraBufferedRateReadExecutor extends AbstractBufferedRateExecu @Value("${cassandra.query.print_queries_freq:0}") int printQueriesFreq, @Autowired StatsFactory statsFactory, @Autowired EntityService entityService, - @Autowired RateLimitService rateLimitService) { + @Autowired RateLimitService rateLimitService, + @Autowired TbServiceInfoProvider serviceInfoProvider) { super(queueLimit, concurrencyLimit, maxWaitTime, dispatcherThreads, callbackThreads, pollMs, printQueriesFreq, statsFactory, - entityService, rateLimitService, printTenantNames); + entityService, rateLimitService, serviceInfoProvider, printTenantNames); } @Scheduled(fixedDelayString = "${cassandra.query.rate_limit_print_interval_ms}") @@ -68,8 +69,8 @@ public class CassandraBufferedRateReadExecutor extends AbstractBufferedRateExecu } @Override - public String getBufferName() { - return BUFFER_NAME; + protected BufferedRateExecutorType getBufferedRateExecutorType() { + return BufferedRateExecutorType.READ; } @Override diff --git a/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraBufferedRateWriteExecutor.java b/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraBufferedRateWriteExecutor.java index 9c894a7eac..c61148e51f 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraBufferedRateWriteExecutor.java +++ b/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraBufferedRateWriteExecutor.java @@ -28,7 +28,9 @@ import org.thingsboard.server.common.stats.StatsFactory; import org.thingsboard.server.dao.entity.EntityService; import org.thingsboard.server.dao.util.AbstractBufferedRateExecutor; import org.thingsboard.server.dao.util.AsyncTaskContext; +import org.thingsboard.server.dao.util.BufferedRateExecutorType; import org.thingsboard.server.dao.util.NoSqlAnyDao; +import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; /** * Created by ashvayka on 24.10.18. @@ -38,8 +40,6 @@ import org.thingsboard.server.dao.util.NoSqlAnyDao; @NoSqlAnyDao public class CassandraBufferedRateWriteExecutor extends AbstractBufferedRateExecutor { - static final String BUFFER_NAME = "Write"; - public CassandraBufferedRateWriteExecutor( @Value("${cassandra.query.buffer_size}") int queueLimit, @Value("${cassandra.query.concurrent_limit}") int concurrencyLimit, @@ -51,9 +51,10 @@ public class CassandraBufferedRateWriteExecutor extends AbstractBufferedRateExec @Value("${cassandra.query.print_queries_freq:0}") int printQueriesFreq, @Autowired StatsFactory statsFactory, @Autowired EntityService entityService, - @Autowired RateLimitService rateLimitService) { + @Autowired RateLimitService rateLimitService, + @Autowired TbServiceInfoProvider serviceInfoProvider) { super(queueLimit, concurrencyLimit, maxWaitTime, dispatcherThreads, callbackThreads, pollMs, printQueriesFreq, statsFactory, - entityService, rateLimitService, printTenantNames); + entityService, rateLimitService, serviceInfoProvider, printTenantNames); } @Scheduled(fixedDelayString = "${cassandra.query.rate_limit_print_interval_ms}") @@ -68,8 +69,8 @@ public class CassandraBufferedRateWriteExecutor extends AbstractBufferedRateExec } @Override - public String getBufferName() { - return BUFFER_NAME; + protected BufferedRateExecutorType getBufferedRateExecutorType() { + return BufferedRateExecutorType.WRITE; } @Override diff --git a/dao/src/main/java/org/thingsboard/server/dao/util/AbstractBufferedRateExecutor.java b/dao/src/main/java/org/thingsboard/server/dao/util/AbstractBufferedRateExecutor.java index a87e9312e3..0f6170f25d 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/util/AbstractBufferedRateExecutor.java +++ b/dao/src/main/java/org/thingsboard/server/dao/util/AbstractBufferedRateExecutor.java @@ -34,12 +34,14 @@ import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.server.cache.limits.RateLimitService; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.limit.LimitedApi; +import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.stats.DefaultCounter; import org.thingsboard.server.common.stats.StatsCounter; import org.thingsboard.server.common.stats.StatsFactory; import org.thingsboard.server.common.stats.StatsType; import org.thingsboard.server.dao.entity.EntityService; import org.thingsboard.server.dao.nosql.CassandraStatementTask; +import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; import java.util.HashMap; import java.util.Map; @@ -78,13 +80,14 @@ public abstract class AbstractBufferedRateExecutor tenantNamesCache = new HashMap<>(); public AbstractBufferedRateExecutor(int queueLimit, int concurrencyLimit, long maxWaitTime, int dispatcherThreads, int callbackThreads, long pollMs, int printQueriesFreq, StatsFactory statsFactory, - EntityService entityService, RateLimitService rateLimitService, boolean printTenantNames) { + EntityService entityService, RateLimitService rateLimitService, TbServiceInfoProvider serviceInfoProvider, boolean printTenantNames) { this.maxWaitTime = maxWaitTime; this.pollMs = pollMs; this.concurrencyLimit = concurrencyLimit; @@ -99,6 +102,7 @@ public abstract class AbstractBufferedRateExecutor execute(AsyncTaskContext taskCtx); - public abstract String getBufferName(); + private String getBufferName() { + return getBufferedRateExecutorType().getDisplayName(); + } + + protected abstract BufferedRateExecutorType getBufferedRateExecutorType(); private void dispatch() { log.info("[{}] Buffered rate executor thread started", getBufferName()); diff --git a/dao/src/main/java/org/thingsboard/server/dao/util/BufferedRateExecutorType.java b/dao/src/main/java/org/thingsboard/server/dao/util/BufferedRateExecutorType.java new file mode 100644 index 0000000000..bf76235e97 --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/util/BufferedRateExecutorType.java @@ -0,0 +1,40 @@ +/** + * Copyright © 2016-2025 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.dao.util; + +import lombok.Getter; +import org.apache.commons.lang3.StringUtils; +import org.thingsboard.server.common.data.limit.LimitedApi; + +@Getter +public enum BufferedRateExecutorType { + + READ(LimitedApi.CASSANDRA_READ_QUERIES_CORE, LimitedApi.CASSANDRA_READ_QUERIES_RULE_ENGINE, LimitedApi.CASSANDRA_READ_QUERIES_MONOLITH), + WRITE(LimitedApi.CASSANDRA_WRITE_QUERIES_CORE, LimitedApi.CASSANDRA_WRITE_QUERIES_RULE_ENGINE, LimitedApi.CASSANDRA_WRITE_QUERIES_MONOLITH); + + private final LimitedApi coreLimitedApi; + private final LimitedApi ruleEngineLimitedApi; + private final LimitedApi monolithLimitedApi; + + private final String displayName = StringUtils.capitalize(name().toLowerCase()); + + BufferedRateExecutorType(LimitedApi coreLimitedApi, LimitedApi ruleEngineLimitedApi, LimitedApi monolithLimitedApi) { + this.coreLimitedApi = coreLimitedApi; + this.ruleEngineLimitedApi = ruleEngineLimitedApi; + this.monolithLimitedApi = monolithLimitedApi; + } + +} From a242d3a43b98ffb4124ff52190d3220501ba66b0 Mon Sep 17 00:00:00 2001 From: dshvaika Date: Tue, 20 May 2025 16:31:01 +0300 Subject: [PATCH 02/10] Added missing test for RateLimitServiceTest && refactoring --- .../service/limits/RateLimitServiceTest.java | 7 ++ .../common/data/limit/LimitedApiEntry.java | 4 - .../common/msg/tools/RateLimitsTest.java | 86 ------------------- .../common/msg/tools/TbRateLimitsTest.java | 66 +++++++++++++- 4 files changed, 71 insertions(+), 92 deletions(-) delete mode 100644 common/message/src/test/java/org/thingsboard/server/common/msg/tools/RateLimitsTest.java diff --git a/application/src/test/java/org/thingsboard/server/service/limits/RateLimitServiceTest.java b/application/src/test/java/org/thingsboard/server/service/limits/RateLimitServiceTest.java index f85ad848ab..674b7bfdb9 100644 --- a/application/src/test/java/org/thingsboard/server/service/limits/RateLimitServiceTest.java +++ b/application/src/test/java/org/thingsboard/server/service/limits/RateLimitServiceTest.java @@ -94,6 +94,13 @@ public class RateLimitServiceTest { testRateLimits(limitedApi, max, tenantId); } + for (LimitedApi limitedApi : List.of( + LimitedApi.CASSANDRA_READ_QUERIES_MONOLITH, + LimitedApi.CASSANDRA_WRITE_QUERIES_MONOLITH + )) { + testRateLimits(limitedApi, max * 2, tenantId); + } + CustomerId customerId = new CustomerId(UUID.randomUUID()); testRateLimits(LimitedApi.REST_REQUESTS_PER_CUSTOMER, max, customerId); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApiEntry.java b/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApiEntry.java index 3082b521b6..16084b7af9 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApiEntry.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApiEntry.java @@ -22,10 +22,6 @@ public record LimitedApiEntry(long capacity, long durationSeconds) { return new LimitedApiEntry(Long.parseLong(parts[0]), Long.parseLong(parts[1])); } - public double rps() { - return (double) capacity / durationSeconds; - } - @Override public String toString() { return capacity + ":" + durationSeconds; diff --git a/common/message/src/test/java/org/thingsboard/server/common/msg/tools/RateLimitsTest.java b/common/message/src/test/java/org/thingsboard/server/common/msg/tools/RateLimitsTest.java deleted file mode 100644 index 9d27687cd9..0000000000 --- a/common/message/src/test/java/org/thingsboard/server/common/msg/tools/RateLimitsTest.java +++ /dev/null @@ -1,86 +0,0 @@ -/** - * Copyright © 2016-2025 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.msg.tools; - -import org.awaitility.pollinterval.FixedPollInterval; -import org.junit.jupiter.api.Test; - -import java.util.concurrent.TimeUnit; - -import static org.assertj.core.api.Assertions.assertThat; -import static org.awaitility.Awaitility.await; - -public class RateLimitsTest { - - @Test - public void testRateLimits_greedyRefill() { - testRateLimitWithGreedyRefill(3, 10); - testRateLimitWithGreedyRefill(3, 3); - testRateLimitWithGreedyRefill(4, 2); - } - - private void testRateLimitWithGreedyRefill(int capacity, int period) { - String rateLimitConfig = capacity + ":" + period; - TbRateLimits rateLimits = new TbRateLimits(rateLimitConfig); - - rateLimits.tryConsume(capacity); - assertThat(rateLimits.tryConsume()).as("new token is available").isFalse(); - - int expectedRefillTime = (int) (((double) period / capacity) * 1000); - int gap = 500; - - for (int i = 0; i < capacity; i++) { - await("token refill for rate limit " + rateLimitConfig) - .pollInterval(new FixedPollInterval(10, TimeUnit.MILLISECONDS)) - .atLeast(expectedRefillTime - gap, TimeUnit.MILLISECONDS) - .atMost(expectedRefillTime + gap, TimeUnit.MILLISECONDS) - .untilAsserted(() -> { - assertThat(rateLimits.tryConsume()).as("token is available").isTrue(); - }); - assertThat(rateLimits.tryConsume()).as("new token is available").isFalse(); - } - } - - @Test - public void testRateLimits_intervalRefill() { - testRateLimitWithIntervalRefill(10, 5); - testRateLimitWithIntervalRefill(3, 3); - testRateLimitWithIntervalRefill(4, 2); - } - - private void testRateLimitWithIntervalRefill(int capacity, int period) { - String rateLimitConfig = capacity + ":" + period; - TbRateLimits rateLimits = new TbRateLimits(rateLimitConfig, true); - - rateLimits.tryConsume(capacity); - assertThat(rateLimits.tryConsume()).as("new token is available").isFalse(); - - int expectedRefillTime = period * 1000; - int gap = 500; - - await("tokens refill for rate limit " + rateLimitConfig) - .pollInterval(new FixedPollInterval(10, TimeUnit.MILLISECONDS)) - .atLeast(expectedRefillTime - gap, TimeUnit.MILLISECONDS) - .atMost(expectedRefillTime + gap, TimeUnit.MILLISECONDS) - .untilAsserted(() -> { - for (int i = 0; i < capacity; i++) { - assertThat(rateLimits.tryConsume()).as("token is available").isTrue(); - } - assertThat(rateLimits.tryConsume()).as("new token is available").isFalse(); - }); - } - -} diff --git a/common/message/src/test/java/org/thingsboard/server/common/msg/tools/TbRateLimitsTest.java b/common/message/src/test/java/org/thingsboard/server/common/msg/tools/TbRateLimitsTest.java index a6a95da9c1..27a2fe6286 100644 --- a/common/message/src/test/java/org/thingsboard/server/common/msg/tools/TbRateLimitsTest.java +++ b/common/message/src/test/java/org/thingsboard/server/common/msg/tools/TbRateLimitsTest.java @@ -15,13 +15,75 @@ */ package org.thingsboard.server.common.msg.tools; +import org.awaitility.pollinterval.FixedPollInterval; import org.junit.jupiter.api.DisplayName; import org.junit.jupiter.api.Test; +import java.util.concurrent.TimeUnit; + import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.awaitility.Awaitility.await; + +public class TbRateLimitsTest { + + @Test + public void testRateLimits_greedyRefill() { + testRateLimitWithGreedyRefill(3, 10); + testRateLimitWithGreedyRefill(3, 3); + testRateLimitWithGreedyRefill(4, 2); + } -class TbRateLimitsTest { + private void testRateLimitWithGreedyRefill(int capacity, int period) { + String rateLimitConfig = capacity + ":" + period; + TbRateLimits rateLimits = new TbRateLimits(rateLimitConfig); + + rateLimits.tryConsume(capacity); + assertThat(rateLimits.tryConsume()).as("new token is available").isFalse(); + + int expectedRefillTime = (int) (((double) period / capacity) * 1000); + int gap = 500; + + for (int i = 0; i < capacity; i++) { + await("token refill for rate limit " + rateLimitConfig) + .pollInterval(new FixedPollInterval(10, TimeUnit.MILLISECONDS)) + .atLeast(expectedRefillTime - gap, TimeUnit.MILLISECONDS) + .atMost(expectedRefillTime + gap, TimeUnit.MILLISECONDS) + .untilAsserted(() -> { + assertThat(rateLimits.tryConsume()).as("token is available").isTrue(); + }); + assertThat(rateLimits.tryConsume()).as("new token is available").isFalse(); + } + } + + @Test + public void testRateLimits_intervalRefill() { + testRateLimitWithIntervalRefill(10, 5); + testRateLimitWithIntervalRefill(3, 3); + testRateLimitWithIntervalRefill(4, 2); + } + + private void testRateLimitWithIntervalRefill(int capacity, int period) { + String rateLimitConfig = capacity + ":" + period; + TbRateLimits rateLimits = new TbRateLimits(rateLimitConfig, true); + + rateLimits.tryConsume(capacity); + assertThat(rateLimits.tryConsume()).as("new token is available").isFalse(); + + int expectedRefillTime = period * 1000; + int gap = 500; + + await("tokens refill for rate limit " + rateLimitConfig) + .pollInterval(new FixedPollInterval(10, TimeUnit.MILLISECONDS)) + .atLeast(expectedRefillTime - gap, TimeUnit.MILLISECONDS) + .atMost(expectedRefillTime + gap, TimeUnit.MILLISECONDS) + .untilAsserted(() -> { + for (int i = 0; i < capacity; i++) { + assertThat(rateLimits.tryConsume()).as("token is available").isTrue(); + } + assertThat(rateLimits.tryConsume()).as("new token is available").isFalse(); + }); + } @Test @DisplayName("TbRateLimits should construct with single rate limit") @@ -60,4 +122,4 @@ class TbRateLimitsTest { .isInstanceOf(ArrayIndexOutOfBoundsException.class); } -} \ No newline at end of file +} From 35e9d007d021914ff2a4ec078b6a3302b06bf80c Mon Sep 17 00:00:00 2001 From: dshvaika Date: Thu, 22 May 2025 18:08:15 +0300 Subject: [PATCH 03/10] upgrade logic + rate limits validation --- .../main/data/upgrade/basic/schema_update.sql | 25 +++++++ .../update/DefaultDataUpdateService.java | 46 +++++++++++++ .../server/common/data/limit/LimitedApi.java | 21 +++--- .../common/data/limit/LimitedApiUtil.java | 25 +++++++ .../DefaultTenantProfileConfiguration.java | 68 +++++++++++++++++++ .../common/data/validation/RateLimit.java | 39 +++++++++++ .../server/common/msg/tools/TbRateLimits.java | 1 - .../dao/service/ConstraintValidator.java | 2 + .../dao/service/RateLimitValidator.java | 32 +++++++++ .../validator/TenantProfileDataValidator.java | 6 +- 10 files changed, 254 insertions(+), 11 deletions(-) create mode 100644 common/data/src/main/java/org/thingsboard/server/common/data/validation/RateLimit.java create mode 100644 dao/src/main/java/org/thingsboard/server/dao/service/RateLimitValidator.java diff --git a/application/src/main/data/upgrade/basic/schema_update.sql b/application/src/main/data/upgrade/basic/schema_update.sql index 016e786776..b19f83621d 100644 --- a/application/src/main/data/upgrade/basic/schema_update.sql +++ b/application/src/main/data/upgrade/basic/schema_update.sql @@ -14,3 +14,28 @@ -- limitations under the License. -- +UPDATE tenant_profile +SET profile_data = jsonb_set( + profile_data, + '{configuration}', + ( + (profile_data -> 'configuration') - 'cassandraQueryTenantRateLimitsConfiguration' + || + COALESCE( + CASE + WHEN profile_data -> 'configuration' -> + 'cassandraQueryTenantRateLimitsConfiguration' IS NOT NULL THEN + jsonb_build_object( + 'cassandraReadQueryTenantCoreRateLimits', + profile_data -> 'configuration' -> 'cassandraQueryTenantRateLimitsConfiguration', + 'cassandraWriteQueryTenantCoreRateLimits', + profile_data -> 'configuration' -> 'cassandraQueryTenantRateLimitsConfiguration', + 'cassandraReadQueryTenantRuleEngineRateLimits', + profile_data -> 'configuration' -> 'cassandraQueryTenantRateLimitsConfiguration', + 'cassandraWriteQueryTenantRuleEngineRateLimits', + profile_data -> 'configuration' -> 'cassandraQueryTenantRateLimitsConfiguration' + ) + END, + '{}'::jsonb) + ) + ); diff --git a/application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java b/application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java index d8ddd7cd9f..c3f1cee046 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java +++ b/application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java @@ -24,6 +24,7 @@ import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.annotation.Profile; import org.springframework.stereotype.Service; +import org.thingsboard.server.common.data.TenantProfile; import org.thingsboard.server.common.data.alarm.AlarmSeverity; import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.RuleNodeId; @@ -34,8 +35,10 @@ import org.thingsboard.server.common.data.query.FilterPredicateValue; import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.data.relation.RelationTypeGroup; import org.thingsboard.server.common.data.rule.RuleNode; +import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration; import org.thingsboard.server.dao.relation.RelationService; import org.thingsboard.server.dao.rule.RuleChainService; +import org.thingsboard.server.dao.tenant.TenantProfileService; import org.thingsboard.server.service.component.ComponentDiscoveryService; import org.thingsboard.server.service.component.RuleNodeClassInfo; import org.thingsboard.server.service.install.DbUpgradeExecutorService; @@ -69,14 +72,57 @@ public class DefaultDataUpdateService implements DataUpdateService { @Autowired private DbUpgradeExecutorService executorService; + @Autowired + private TenantProfileService tenantProfileService; + @Override public void updateData() throws Exception { log.info("Updating data ..."); //TODO: should be cleaned after each release updateInputNodes(); + deduplicateRateLimitsPerSecondsConfigurations(); log.info("Data updated."); } + private void deduplicateRateLimitsPerSecondsConfigurations() { + log.info("Starting update of tenant profiles..."); + + int totalProfiles = 0; + int updatedTenantProfiles = 0; + int skippedProfiles = 0; + int failedProfiles = 0; + + var tenantProfiles = new PageDataIterable<>( + pageLink -> tenantProfileService.findTenantProfiles(TenantId.SYS_TENANT_ID, pageLink), 1024); + + for (TenantProfile tenantProfile : tenantProfiles) { + totalProfiles++; + String profileName = tenantProfile.getName(); + UUID profileId = tenantProfile.getId().getId(); + try { + Optional profileConfiguration = tenantProfile.getProfileConfiguration(); + if (profileConfiguration.isEmpty()) { + log.debug("[{}][{}] Skipping tenant profile with non-default configuration.", profileId, profileName); + skippedProfiles++; + continue; + } + + DefaultTenantProfileConfiguration defaultTenantProfileConfiguration = profileConfiguration.get(); + defaultTenantProfileConfiguration.deduplicateRateLimitsConfigs(); + tenantProfileService.saveTenantProfile(TenantId.SYS_TENANT_ID, tenantProfile); + updatedTenantProfiles++; + log.debug("[{}][{}] Successfully updated tenant profile.", profileId, profileName); + } catch (Exception e) { + log.error("[{}][{}] Failed to updated tenant profile: ", profileId, profileName, e); + failedProfiles++; + } + } + + log.info("Tenant profiles update completed. Total: {}, Updated: {}, Skipped: {}, Failed: {}", + totalProfiles, updatedTenantProfiles, skippedProfiles, failedProfiles); + } + + private void updateInputNodes() { log.info("Creating relations for input nodes..."); int n = 0; diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApi.java b/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApi.java index 766a1f4eb3..428c430c86 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApi.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApi.java @@ -31,7 +31,7 @@ public enum LimitedApi { REST_REQUESTS_PER_CUSTOMER(DefaultTenantProfileConfiguration::getCustomerServerRestLimitsConfiguration, "REST API requests per customer", false), WS_UPDATES_PER_SESSION(DefaultTenantProfileConfiguration::getWsUpdatesPerSessionRateLimit, "WS updates per session", true), CASSANDRA_WRITE_QUERIES_CORE(DefaultTenantProfileConfiguration::getCassandraReadQueryTenantCoreRateLimits, "Rest API and WS telemetry read queries", true), - CASSANDRA_READ_QUERIES_CORE(DefaultTenantProfileConfiguration::getCassandraWriteQueryTenantCoreRateLimits, "Rest API and WS telemetry write queries", true), + CASSANDRA_READ_QUERIES_CORE(DefaultTenantProfileConfiguration::getCassandraWriteQueryTenantCoreRateLimits, "Rest API write queries", true), CASSANDRA_WRITE_QUERIES_RULE_ENGINE(DefaultTenantProfileConfiguration::getCassandraReadQueryTenantRuleEngineRateLimits, "Rule Engine telemetry read queries", true), CASSANDRA_READ_QUERIES_RULE_ENGINE(DefaultTenantProfileConfiguration::getCassandraWriteQueryTenantRuleEngineRateLimits, "Rule Engine telemetry write queries", true), CASSANDRA_READ_QUERIES_MONOLITH( @@ -57,28 +57,31 @@ public enum LimitedApi { WS_SUBSCRIPTIONS("WS subscriptions", false), CALCULATED_FIELD_DEBUG_EVENTS("calculated field debug events", true); - private Function configExtractor; + private final Function configExtractor; @Getter private final boolean perTenant; @Getter - private boolean refillRateLimitIntervally; + private final boolean refillRateLimitIntervally; @Getter - private String label; + private final String label; LimitedApi(Function configExtractor, String label, boolean perTenant) { - this.configExtractor = configExtractor; - this.label = label; - this.perTenant = perTenant; + this(configExtractor, label, perTenant, false); } LimitedApi(boolean perTenant, boolean refillRateLimitIntervally) { - this.perTenant = perTenant; - this.refillRateLimitIntervally = refillRateLimitIntervally; + this(null, null, perTenant, refillRateLimitIntervally); } LimitedApi(String label, boolean perTenant) { + this(null, label, perTenant, false); + } + + LimitedApi(Function configExtractor, String label, boolean perTenant, boolean refillRateLimitIntervally) { + this.configExtractor = configExtractor; this.label = label; this.perTenant = perTenant; + this.refillRateLimitIntervally = refillRateLimitIntervally; } public String getLimitConfig(DefaultTenantProfileConfiguration profileConfiguration) { diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApiUtil.java b/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApiUtil.java index a90c7eb825..62a8d6567f 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApiUtil.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApiUtil.java @@ -21,8 +21,10 @@ import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; import java.util.HashMap; +import java.util.HashSet; import java.util.List; import java.util.Map; +import java.util.Set; import java.util.function.Function; import java.util.stream.Collectors; @@ -64,4 +66,27 @@ public class LimitedApiUtil { .collect(Collectors.joining(",")); } + public static boolean isValid(String configStr) { + List limitedApiEntries = parseConfig(configStr); + Set distinctDurations = new HashSet<>(); + for (LimitedApiEntry entry : limitedApiEntries) { + if (!distinctDurations.add(entry.durationSeconds())) { + return false; + } + } + return true; + } + + @Deprecated(forRemoval = true, since = "4.0.2") + public static String deduplicateByDuration(String configStr) { + if (configStr == null) { + return null; + } + Set distinctDurations = new HashSet<>(); + return parseConfig(configStr).stream() + .filter(entry -> distinctDurations.add(entry.durationSeconds())) + .map(LimitedApiEntry::toString) + .collect(Collectors.joining(",")); + } + } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java index d61f48b148..4522fb75bf 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java @@ -24,6 +24,8 @@ import lombok.NoArgsConstructor; import org.thingsboard.server.common.data.ApiUsageRecordKey; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.TenantProfileType; +import org.thingsboard.server.common.data.limit.LimitedApiUtil; +import org.thingsboard.server.common.data.validation.RateLimit; import java.io.Serial; @@ -49,37 +51,53 @@ public class DefaultTenantProfileConfiguration implements TenantProfileConfigura private long maxResourceSize; @Schema(example = "1000:1,20000:60") + @RateLimit(fieldName = "Transport tenant messages") private String transportTenantMsgRateLimit; @Schema(example = "1000:1,20000:60") + @RateLimit(fieldName = "Transport tenant telemetry messages") private String transportTenantTelemetryMsgRateLimit; @Schema(example = "1000:1,20000:60") + @RateLimit(fieldName = "Transport tenant telemetry data points") private String transportTenantTelemetryDataPointsRateLimit; @Schema(example = "20:1,600:60") + @RateLimit(fieldName = "Transport device messages") private String transportDeviceMsgRateLimit; @Schema(example = "20:1,600:60") + @RateLimit(fieldName = "Transport device telemetry messages") private String transportDeviceTelemetryMsgRateLimit; @Schema(example = "20:1,600:60") + @RateLimit(fieldName = "Transport device telemetry data points") private String transportDeviceTelemetryDataPointsRateLimit; @Schema(example = "20:1,600:60") + @RateLimit(fieldName = "Transport gateway messages") private String transportGatewayMsgRateLimit; @Schema(example = "20:1,600:60") + @RateLimit(fieldName = "Transport gateway telemetry messages") private String transportGatewayTelemetryMsgRateLimit; @Schema(example = "20:1,600:60") + @RateLimit(fieldName = "Transport gateway telemetry data points") private String transportGatewayTelemetryDataPointsRateLimit; @Schema(example = "20:1,600:60") + @RateLimit(fieldName = "Transport gateway device messages") private String transportGatewayDeviceMsgRateLimit; @Schema(example = "20:1,600:60") + @RateLimit(fieldName = "Transport gateway device telemetry messages") private String transportGatewayDeviceTelemetryMsgRateLimit; @Schema(example = "20:1,600:60") + @RateLimit(fieldName = "Transport gateway device telemetry data points") private String transportGatewayDeviceTelemetryDataPointsRateLimit; @Schema(example = "20:1,600:60") + @RateLimit(fieldName = "Entity version creation") private String tenantEntityExportRateLimit; @Schema(example = "20:1,600:60") + @RateLimit(fieldName = "Entity version load") private String tenantEntityImportRateLimit; @Schema(example = "20:1,600:60") + @RateLimit(fieldName = "Notification requests") private String tenantNotificationRequestsRateLimit; @Schema(example = "20:1,600:60") + @RateLimit(fieldName = "Notification requests per notification rule") private String tenantNotificationRequestsPerRuleRateLimit; @Schema(example = "10000000") @@ -107,7 +125,9 @@ public class DefaultTenantProfileConfiguration implements TenantProfileConfigura @Schema(example = "1000") private long maxCreatedAlarms; + @RateLimit(fieldName = "REST requests for tenant") private String tenantServerRestLimitsConfiguration; + @RateLimit(fieldName = "REST requests for customer") private String customerServerRestLimitsConfiguration; private int maxWsSessionsPerTenant; @@ -119,17 +139,26 @@ public class DefaultTenantProfileConfiguration implements TenantProfileConfigura private long maxWsSubscriptionsPerCustomer; private long maxWsSubscriptionsPerRegularUser; private long maxWsSubscriptionsPerPublicUser; + @RateLimit(fieldName = "WS updates per session") private String wsUpdatesPerSessionRateLimit; + @RateLimit(fieldName = "Rest API and WS telemetry read queries") private String cassandraReadQueryTenantCoreRateLimits; + @RateLimit(fieldName = "Rest API write queries") private String cassandraWriteQueryTenantCoreRateLimits; + @RateLimit(fieldName = "Rule Engine telemetry read queries") private String cassandraReadQueryTenantRuleEngineRateLimits; + @RateLimit(fieldName = "Rule Engine telemetry write queries") private String cassandraWriteQueryTenantRuleEngineRateLimits; + @RateLimit(fieldName = "Edge events") private String edgeEventRateLimits; + @RateLimit(fieldName = "Edge events per edge") private String edgeEventRateLimitsPerEdge; + @RateLimit(fieldName = "Edge uplink messages") private String edgeUplinkMessagesRateLimits; + @RateLimit(fieldName = "Edge uplink messages per edge") private String edgeUplinkMessagesRateLimitsPerEdge; private int defaultStorageTtlDays; @@ -207,4 +236,43 @@ public class DefaultTenantProfileConfiguration implements TenantProfileConfigura return maxRuleNodeExecutionsPerMessage; } + @Deprecated(forRemoval = true, since = "4.0.2") + public void deduplicateRateLimitsConfigs() { + this.transportTenantMsgRateLimit = LimitedApiUtil.deduplicateByDuration(transportTenantMsgRateLimit); + this.transportTenantTelemetryMsgRateLimit = LimitedApiUtil.deduplicateByDuration(transportTenantTelemetryMsgRateLimit); + this.transportTenantTelemetryDataPointsRateLimit = LimitedApiUtil.deduplicateByDuration(transportTenantTelemetryDataPointsRateLimit); + + this.transportDeviceMsgRateLimit = LimitedApiUtil.deduplicateByDuration(transportDeviceMsgRateLimit); + this.transportDeviceTelemetryMsgRateLimit = LimitedApiUtil.deduplicateByDuration(transportDeviceTelemetryMsgRateLimit); + this.transportDeviceTelemetryDataPointsRateLimit = LimitedApiUtil.deduplicateByDuration(transportDeviceTelemetryDataPointsRateLimit); + + this.transportGatewayMsgRateLimit = LimitedApiUtil.deduplicateByDuration(transportGatewayMsgRateLimit); + this.transportGatewayTelemetryMsgRateLimit = LimitedApiUtil.deduplicateByDuration(transportGatewayTelemetryMsgRateLimit); + this.transportGatewayTelemetryDataPointsRateLimit = LimitedApiUtil.deduplicateByDuration(transportGatewayTelemetryDataPointsRateLimit); + + this.transportGatewayDeviceMsgRateLimit = LimitedApiUtil.deduplicateByDuration(transportGatewayDeviceMsgRateLimit); + this.transportGatewayDeviceTelemetryMsgRateLimit = LimitedApiUtil.deduplicateByDuration(transportGatewayDeviceTelemetryMsgRateLimit); + this.transportGatewayDeviceTelemetryDataPointsRateLimit = LimitedApiUtil.deduplicateByDuration(transportGatewayDeviceTelemetryDataPointsRateLimit); + + this.tenantEntityExportRateLimit = LimitedApiUtil.deduplicateByDuration(tenantEntityExportRateLimit); + this.tenantEntityImportRateLimit = LimitedApiUtil.deduplicateByDuration(tenantEntityImportRateLimit); + this.tenantNotificationRequestsRateLimit = LimitedApiUtil.deduplicateByDuration(tenantNotificationRequestsRateLimit); + this.tenantNotificationRequestsPerRuleRateLimit = LimitedApiUtil.deduplicateByDuration(tenantNotificationRequestsPerRuleRateLimit); + + this.cassandraReadQueryTenantCoreRateLimits = LimitedApiUtil.deduplicateByDuration(cassandraReadQueryTenantCoreRateLimits); + this.cassandraWriteQueryTenantCoreRateLimits = LimitedApiUtil.deduplicateByDuration(cassandraWriteQueryTenantCoreRateLimits); + this.cassandraReadQueryTenantRuleEngineRateLimits = LimitedApiUtil.deduplicateByDuration(cassandraReadQueryTenantRuleEngineRateLimits); + this.cassandraWriteQueryTenantRuleEngineRateLimits = LimitedApiUtil.deduplicateByDuration(cassandraWriteQueryTenantRuleEngineRateLimits); + + this.edgeEventRateLimits = LimitedApiUtil.deduplicateByDuration(edgeEventRateLimits); + this.edgeEventRateLimitsPerEdge = LimitedApiUtil.deduplicateByDuration(edgeEventRateLimitsPerEdge); + this.edgeUplinkMessagesRateLimits = LimitedApiUtil.deduplicateByDuration(edgeUplinkMessagesRateLimits); + this.edgeUplinkMessagesRateLimitsPerEdge = LimitedApiUtil.deduplicateByDuration(edgeUplinkMessagesRateLimitsPerEdge); + + this.wsUpdatesPerSessionRateLimit = LimitedApiUtil.deduplicateByDuration(wsUpdatesPerSessionRateLimit); + + this.tenantServerRestLimitsConfiguration = LimitedApiUtil.deduplicateByDuration(tenantServerRestLimitsConfiguration); + this.customerServerRestLimitsConfiguration = LimitedApiUtil.deduplicateByDuration(customerServerRestLimitsConfiguration); + } + } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/validation/RateLimit.java b/common/data/src/main/java/org/thingsboard/server/common/data/validation/RateLimit.java new file mode 100644 index 0000000000..9afe71da2f --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/validation/RateLimit.java @@ -0,0 +1,39 @@ +/** + * Copyright © 2016-2025 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.data.validation; + +import jakarta.validation.Constraint; +import jakarta.validation.Payload; + +import java.lang.annotation.ElementType; +import java.lang.annotation.Retention; +import java.lang.annotation.RetentionPolicy; +import java.lang.annotation.Target; + +@Retention(RetentionPolicy.RUNTIME) +@Target({ElementType.FIELD, ElementType.METHOD}) +@Constraint(validatedBy = {}) +public @interface RateLimit { + + String message() default "rate limit has duplicate 'Per seconds' configuration."; + + String fieldName() default ""; + + Class[] groups() default {}; + + Class[] payload() default {}; + +} diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/tools/TbRateLimits.java b/common/message/src/main/java/org/thingsboard/server/common/msg/tools/TbRateLimits.java index 1e2d8dd4ce..c65c90706a 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/tools/TbRateLimits.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/tools/TbRateLimits.java @@ -18,7 +18,6 @@ package org.thingsboard.server.common.msg.tools; import io.github.bucket4j.Bandwidth; import io.github.bucket4j.BandwidthBuilder; import io.github.bucket4j.Bucket; -import io.github.bucket4j.Refill; import io.github.bucket4j.local.LocalBucket; import io.github.bucket4j.local.LocalBucketBuilder; import lombok.Getter; diff --git a/dao/src/main/java/org/thingsboard/server/dao/service/ConstraintValidator.java b/dao/src/main/java/org/thingsboard/server/dao/service/ConstraintValidator.java index 4992af5b96..f836f229ea 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/service/ConstraintValidator.java +++ b/dao/src/main/java/org/thingsboard/server/dao/service/ConstraintValidator.java @@ -33,6 +33,7 @@ import org.springframework.context.annotation.Configuration; import org.springframework.validation.beanvalidation.LocalValidatorFactoryBean; import org.thingsboard.server.common.data.validation.Length; import org.thingsboard.server.common.data.validation.NoXss; +import org.thingsboard.server.common.data.validation.RateLimit; import org.thingsboard.server.dao.exception.DataValidationException; import java.util.Collection; @@ -103,6 +104,7 @@ public class ConstraintValidator { ConstraintMapping constraintMapping = new DefaultConstraintMapping(null); constraintMapping.constraintDefinition(NoXss.class).validatedBy(NoXssValidator.class); constraintMapping.constraintDefinition(Length.class).validatedBy(StringLengthValidator.class); + constraintMapping.constraintDefinition(RateLimit.class).validatedBy(RateLimitValidator.class); return constraintMapping; } diff --git a/dao/src/main/java/org/thingsboard/server/dao/service/RateLimitValidator.java b/dao/src/main/java/org/thingsboard/server/dao/service/RateLimitValidator.java new file mode 100644 index 0000000000..0703f8b65d --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/service/RateLimitValidator.java @@ -0,0 +1,32 @@ +/** + * Copyright © 2016-2025 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.dao.service; + +import jakarta.validation.ConstraintValidator; +import jakarta.validation.ConstraintValidatorContext; +import lombok.extern.slf4j.Slf4j; +import org.thingsboard.server.common.data.limit.LimitedApiUtil; +import org.thingsboard.server.common.data.validation.RateLimit; + +@Slf4j +public class RateLimitValidator implements ConstraintValidator { + + @Override + public boolean isValid(String value, ConstraintValidatorContext constraintValidatorContext) { + return value == null || LimitedApiUtil.isValid(value); + } + +} diff --git a/dao/src/main/java/org/thingsboard/server/dao/service/validator/TenantProfileDataValidator.java b/dao/src/main/java/org/thingsboard/server/dao/service/validator/TenantProfileDataValidator.java index fae3ab43d6..6f6b6a3528 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/service/validator/TenantProfileDataValidator.java +++ b/dao/src/main/java/org/thingsboard/server/dao/service/validator/TenantProfileDataValidator.java @@ -24,6 +24,7 @@ import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.queue.ProcessingStrategy; import org.thingsboard.server.common.data.queue.SubmitStrategy; import org.thingsboard.server.common.data.queue.SubmitStrategyType; +import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration; import org.thingsboard.server.common.data.tenant.profile.TenantProfileQueueConfiguration; import org.thingsboard.server.dao.exception.DataValidationException; import org.thingsboard.server.dao.service.DataValidator; @@ -51,9 +52,12 @@ public class TenantProfileDataValidator extends DataValidator { if (tenantProfile.getProfileData() == null) { throw new DataValidationException("Tenant profile data should be specified!"); } - if (tenantProfile.getProfileData().getConfiguration() == null) { + + Optional profileConfiguration = tenantProfile.getProfileConfiguration(); + if (profileConfiguration.isEmpty()) { throw new DataValidationException("Tenant profile data configuration should be specified!"); } + if (tenantProfile.isDefault()) { TenantProfile defaultTenantProfile = tenantProfileService.findDefaultTenantProfile(tenantId); if (defaultTenantProfile != null && !defaultTenantProfile.getId().equals(tenantProfile.getId())) { From 508fa081bff6274623054e89937d3fd23014a13c Mon Sep 17 00:00:00 2001 From: dshvaika Date: Fri, 23 May 2025 13:32:35 +0300 Subject: [PATCH 04/10] made serviceInfoProvider not required for install application --- .../server/dao/nosql/CassandraBufferedRateReadExecutor.java | 2 +- .../server/dao/nosql/CassandraBufferedRateWriteExecutor.java | 2 +- .../server/dao/util/AbstractBufferedRateExecutor.java | 3 +++ 3 files changed, 5 insertions(+), 2 deletions(-) diff --git a/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraBufferedRateReadExecutor.java b/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraBufferedRateReadExecutor.java index a7aa0dc5dc..98cfef7f7d 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraBufferedRateReadExecutor.java +++ b/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraBufferedRateReadExecutor.java @@ -52,7 +52,7 @@ public class CassandraBufferedRateReadExecutor extends AbstractBufferedRateExecu @Autowired StatsFactory statsFactory, @Autowired EntityService entityService, @Autowired RateLimitService rateLimitService, - @Autowired TbServiceInfoProvider serviceInfoProvider) { + @Autowired(required = false) TbServiceInfoProvider serviceInfoProvider) { super(queueLimit, concurrencyLimit, maxWaitTime, dispatcherThreads, callbackThreads, pollMs, printQueriesFreq, statsFactory, entityService, rateLimitService, serviceInfoProvider, printTenantNames); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraBufferedRateWriteExecutor.java b/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraBufferedRateWriteExecutor.java index c61148e51f..1fb4b8dea5 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraBufferedRateWriteExecutor.java +++ b/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraBufferedRateWriteExecutor.java @@ -52,7 +52,7 @@ public class CassandraBufferedRateWriteExecutor extends AbstractBufferedRateExec @Autowired StatsFactory statsFactory, @Autowired EntityService entityService, @Autowired RateLimitService rateLimitService, - @Autowired TbServiceInfoProvider serviceInfoProvider) { + @Autowired(required = false) TbServiceInfoProvider serviceInfoProvider) { super(queueLimit, concurrencyLimit, maxWaitTime, dispatcherThreads, callbackThreads, pollMs, printQueriesFreq, statsFactory, entityService, rateLimitService, serviceInfoProvider, printTenantNames); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/util/AbstractBufferedRateExecutor.java b/dao/src/main/java/org/thingsboard/server/dao/util/AbstractBufferedRateExecutor.java index 0f6170f25d..444aecef71 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/util/AbstractBufferedRateExecutor.java +++ b/dao/src/main/java/org/thingsboard/server/dao/util/AbstractBufferedRateExecutor.java @@ -141,6 +141,9 @@ public abstract class AbstractBufferedRateExecutor Date: Fri, 23 May 2025 14:14:42 +0300 Subject: [PATCH 05/10] minor improvement to upgrade script logic --- application/src/main/data/upgrade/basic/schema_update.sql | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/application/src/main/data/upgrade/basic/schema_update.sql b/application/src/main/data/upgrade/basic/schema_update.sql index b19f83621d..ecb2a147bc 100644 --- a/application/src/main/data/upgrade/basic/schema_update.sql +++ b/application/src/main/data/upgrade/basic/schema_update.sql @@ -36,6 +36,9 @@ SET profile_data = jsonb_set( profile_data -> 'configuration' -> 'cassandraQueryTenantRateLimitsConfiguration' ) END, - '{}'::jsonb) + '{}'::jsonb ) - ); + ) + ) +WHERE profile_data -> 'configuration' ? 'cassandraQueryTenantRateLimitsConfiguration'; + From 6348035617113b4a7a0b2977172473839c0cdb78 Mon Sep 17 00:00:00 2001 From: dshvaika Date: Fri, 23 May 2025 15:38:26 +0300 Subject: [PATCH 06/10] UI part + minor refactoring --- .../server/common/data/limit/LimitedApi.java | 8 ++++---- .../common/data/limit/LimitedApiUtil.java | 2 +- .../DefaultTenantProfileConfiguration.java | 10 +++++----- .../common/data/limit/LimitedApiUtilTest.java | 14 ------------- ...enant-profile-configuration.component.html | 20 +++++++++++++------ ...-tenant-profile-configuration.component.ts | 5 ++++- .../tenant/rate-limits/rate-limits.models.ts | 15 +++++++++++--- ui-ngx/src/app/shared/models/tenant.model.ts | 10 ++++++++-- .../assets/locale/locale.constant-en_US.json | 10 ++++++++-- 9 files changed, 56 insertions(+), 38 deletions(-) diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApi.java b/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApi.java index 428c430c86..19eb4b510f 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApi.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApi.java @@ -30,10 +30,10 @@ public enum LimitedApi { REST_REQUESTS_PER_TENANT(DefaultTenantProfileConfiguration::getTenantServerRestLimitsConfiguration, "REST API requests", true), REST_REQUESTS_PER_CUSTOMER(DefaultTenantProfileConfiguration::getCustomerServerRestLimitsConfiguration, "REST API requests per customer", false), WS_UPDATES_PER_SESSION(DefaultTenantProfileConfiguration::getWsUpdatesPerSessionRateLimit, "WS updates per session", true), - CASSANDRA_WRITE_QUERIES_CORE(DefaultTenantProfileConfiguration::getCassandraReadQueryTenantCoreRateLimits, "Rest API and WS telemetry read queries", true), - CASSANDRA_READ_QUERIES_CORE(DefaultTenantProfileConfiguration::getCassandraWriteQueryTenantCoreRateLimits, "Rest API write queries", true), - CASSANDRA_WRITE_QUERIES_RULE_ENGINE(DefaultTenantProfileConfiguration::getCassandraReadQueryTenantRuleEngineRateLimits, "Rule Engine telemetry read queries", true), - CASSANDRA_READ_QUERIES_RULE_ENGINE(DefaultTenantProfileConfiguration::getCassandraWriteQueryTenantRuleEngineRateLimits, "Rule Engine telemetry write queries", true), + CASSANDRA_WRITE_QUERIES_CORE(DefaultTenantProfileConfiguration::getCassandraReadQueryTenantCoreRateLimits, "Rest API Cassandra write queries", true), + CASSANDRA_READ_QUERIES_CORE(DefaultTenantProfileConfiguration::getCassandraWriteQueryTenantCoreRateLimits, "Rest API and WS telemetry Cassandra read queries", true), + CASSANDRA_WRITE_QUERIES_RULE_ENGINE(DefaultTenantProfileConfiguration::getCassandraReadQueryTenantRuleEngineRateLimits, "Rule Engine telemetry Cassandra write queries", true), + CASSANDRA_READ_QUERIES_RULE_ENGINE(DefaultTenantProfileConfiguration::getCassandraWriteQueryTenantRuleEngineRateLimits, "Rule Engine telemetry Cassandra read queries", true), CASSANDRA_READ_QUERIES_MONOLITH( LimitedApiUtil.merge( DefaultTenantProfileConfiguration::getCassandraReadQueryTenantCoreRateLimits, diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApiUtil.java b/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApiUtil.java index 62a8d6567f..9c99d065d4 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApiUtil.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApiUtil.java @@ -77,7 +77,7 @@ public class LimitedApiUtil { return true; } - @Deprecated(forRemoval = true, since = "4.0.2") + @Deprecated(forRemoval = true, since = "4.1") public static String deduplicateByDuration(String configStr) { if (configStr == null) { return null; diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java index 4522fb75bf..00fecfe800 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java @@ -142,14 +142,14 @@ public class DefaultTenantProfileConfiguration implements TenantProfileConfigura @RateLimit(fieldName = "WS updates per session") private String wsUpdatesPerSessionRateLimit; - @RateLimit(fieldName = "Rest API and WS telemetry read queries") + @RateLimit(fieldName = "Rest API and WS telemetry Cassandra read queries") private String cassandraReadQueryTenantCoreRateLimits; - @RateLimit(fieldName = "Rest API write queries") + @RateLimit(fieldName = "Rest API Cassandra write queries") private String cassandraWriteQueryTenantCoreRateLimits; - @RateLimit(fieldName = "Rule Engine telemetry read queries") + @RateLimit(fieldName = "Rule Engine telemetry Cassandra read queries") private String cassandraReadQueryTenantRuleEngineRateLimits; - @RateLimit(fieldName = "Rule Engine telemetry write queries") + @RateLimit(fieldName = "Rule Engine telemetry Cassandra write queries") private String cassandraWriteQueryTenantRuleEngineRateLimits; @RateLimit(fieldName = "Edge events") @@ -236,7 +236,7 @@ public class DefaultTenantProfileConfiguration implements TenantProfileConfigura return maxRuleNodeExecutionsPerMessage; } - @Deprecated(forRemoval = true, since = "4.0.2") + @Deprecated(forRemoval = true, since = "4.1") public void deduplicateRateLimitsConfigs() { this.transportTenantMsgRateLimit = LimitedApiUtil.deduplicateByDuration(transportTenantMsgRateLimit); this.transportTenantTelemetryMsgRateLimit = LimitedApiUtil.deduplicateByDuration(transportTenantTelemetryMsgRateLimit); diff --git a/common/data/src/test/java/org/thingsboard/server/common/data/limit/LimitedApiUtilTest.java b/common/data/src/test/java/org/thingsboard/server/common/data/limit/LimitedApiUtilTest.java index 6486de3221..c275df790c 100644 --- a/common/data/src/test/java/org/thingsboard/server/common/data/limit/LimitedApiUtilTest.java +++ b/common/data/src/test/java/org/thingsboard/server/common/data/limit/LimitedApiUtilTest.java @@ -96,18 +96,4 @@ class LimitedApiUtilTest { assertThat(result).isEqualTo("200:10,100:60"); } - @Test - @DisplayName("LimitedApiUtil shouldn't have duplicate durations in the same config!") - void testMergeHandlesDuplicatesInSingleConfig() { - Function extractor1 = cfg -> "100:60,200:60"; - Function extractor2 = cfg -> ""; - - // Fake config instance (not used directly in lambda logic) - DefaultTenantProfileConfiguration config = new DefaultTenantProfileConfiguration(); - String result = LimitedApiUtil.merge(extractor1, extractor2).apply(config); - - // 100+200 = 300 for duration 60. Currently possible to save the same "per seconds" config from the UI. - // This must be fixed, so we will merge only two different rate limits. - assertThat(result).isEqualTo("300:60"); - } } diff --git a/ui-ngx/src/app/modules/home/components/profile/tenant/default-tenant-profile-configuration.component.html b/ui-ngx/src/app/modules/home/components/profile/tenant/default-tenant-profile-configuration.component.html index aded3e2f78..3c5c8774d5 100644 --- a/ui-ngx/src/app/modules/home/components/profile/tenant/default-tenant-profile-configuration.component.html +++ b/ui-ngx/src/app/modules/home/components/profile/tenant/default-tenant-profile-configuration.component.html @@ -693,11 +693,19 @@
- + + + + +
+
+ - +
@@ -725,8 +733,8 @@
- +
diff --git a/ui-ngx/src/app/modules/home/components/profile/tenant/default-tenant-profile-configuration.component.ts b/ui-ngx/src/app/modules/home/components/profile/tenant/default-tenant-profile-configuration.component.ts index b1d6652e4a..c6efd9dee3 100644 --- a/ui-ngx/src/app/modules/home/components/profile/tenant/default-tenant-profile-configuration.component.ts +++ b/ui-ngx/src/app/modules/home/components/profile/tenant/default-tenant-profile-configuration.component.ts @@ -114,7 +114,10 @@ export class DefaultTenantProfileConfigurationComponent implements ControlValueA maxWsSubscriptionsPerRegularUser: [null, [Validators.min(0)]], maxWsSubscriptionsPerPublicUser: [null, [Validators.min(0)]], wsUpdatesPerSessionRateLimit: [null, []], - cassandraQueryTenantRateLimitsConfiguration: [null, []], + cassandraWriteQueryTenantCoreRateLimits: [null, []], + cassandraReadQueryTenantCoreRateLimits: [null, []], + cassandraWriteQueryTenantRuleEngineRateLimits: [null, []], + cassandraReadQueryTenantRuleEngineRateLimits: [null, []], edgeEventRateLimits: [null, []], edgeEventRateLimitsPerEdge: [null, []], edgeUplinkMessagesRateLimits: [null, []], diff --git a/ui-ngx/src/app/modules/home/components/profile/tenant/rate-limits/rate-limits.models.ts b/ui-ngx/src/app/modules/home/components/profile/tenant/rate-limits/rate-limits.models.ts index f09f950ee7..bfa8b9ac1d 100644 --- a/ui-ngx/src/app/modules/home/components/profile/tenant/rate-limits/rate-limits.models.ts +++ b/ui-ngx/src/app/modules/home/components/profile/tenant/rate-limits/rate-limits.models.ts @@ -37,7 +37,10 @@ export enum RateLimitsType { TENANT_SERVER_REST_LIMITS_CONFIGURATION = 'TENANT_SERVER_REST_LIMITS_CONFIGURATION', CUSTOMER_SERVER_REST_LIMITS_CONFIGURATION = 'CUSTOMER_SERVER_REST_LIMITS_CONFIGURATION', WS_UPDATE_PER_SESSION_RATE_LIMIT = 'WS_UPDATE_PER_SESSION_RATE_LIMIT', - CASSANDRA_QUERY_TENANT_RATE_LIMITS_CONFIGURATION = 'CASSANDRA_QUERY_TENANT_RATE_LIMITS_CONFIGURATION', + CASSANDRA_WRITE_QUERY_TENANT_CORE_RATE_LIMITS = 'CASSANDRA_WRITE_QUERY_TENANT_CORE_RATE_LIMITS', + CASSANDRA_READ_QUERY_TENANT_CORE_RATE_LIMITS = 'CASSANDRA_READ_QUERY_TENANT_CORE_RATE_LIMITS', + CASSANDRA_WRITE_QUERY_TENANT_RULE_ENGINE_RATE_LIMITS = 'CASSANDRA_WRITE_QUERY_TENANT_RULE_ENGINE_RATE_LIMITS', + CASSANDRA_READ_QUERY_TENANT_RULE_ENGINE_RATE_LIMITS = 'CASSANDRA_READ_QUERY_TENANT_RULE_ENGINE_RATE_LIMITS', TENANT_ENTITY_EXPORT_RATE_LIMIT = 'TENANT_ENTITY_EXPORT_RATE_LIMIT', TENANT_ENTITY_IMPORT_RATE_LIMIT = 'TENANT_ENTITY_IMPORT_RATE_LIMIT', TENANT_NOTIFICATION_REQUEST_RATE_LIMIT = 'TENANT_NOTIFICATION_REQUEST_RATE_LIMIT', @@ -66,7 +69,10 @@ export const rateLimitsLabelTranslationMap = new Map( [RateLimitsType.TENANT_SERVER_REST_LIMITS_CONFIGURATION, 'tenant-profile.rest-requests-for-tenant'], [RateLimitsType.CUSTOMER_SERVER_REST_LIMITS_CONFIGURATION, 'tenant-profile.customer-rest-limits'], [RateLimitsType.WS_UPDATE_PER_SESSION_RATE_LIMIT, 'tenant-profile.ws-limit-updates-per-session'], - [RateLimitsType.CASSANDRA_QUERY_TENANT_RATE_LIMITS_CONFIGURATION, 'tenant-profile.cassandra-tenant-limits-configuration'], + [RateLimitsType.CASSANDRA_WRITE_QUERY_TENANT_CORE_RATE_LIMITS, 'tenant-profile.cassandra-write-tenant-core-limits-configuration'], + [RateLimitsType.CASSANDRA_READ_QUERY_TENANT_CORE_RATE_LIMITS, 'tenant-profile.cassandra-read-tenant-core-limits-configuration'], + [RateLimitsType.CASSANDRA_WRITE_QUERY_TENANT_RULE_ENGINE_RATE_LIMITS, 'tenant-profile.cassandra-write-tenant-rule-engine-limits-configuration'], + [RateLimitsType.CASSANDRA_READ_QUERY_TENANT_RULE_ENGINE_RATE_LIMITS, 'tenant-profile.cassandra-read-tenant-rule-engine-limits-configuration'], [RateLimitsType.TENANT_ENTITY_EXPORT_RATE_LIMIT, 'tenant-profile.tenant-entity-export-rate-limit'], [RateLimitsType.TENANT_ENTITY_IMPORT_RATE_LIMIT, 'tenant-profile.tenant-entity-import-rate-limit'], [RateLimitsType.TENANT_NOTIFICATION_REQUEST_RATE_LIMIT, 'tenant-profile.tenant-notification-request-rate-limit'], @@ -96,7 +102,10 @@ export const rateLimitsDialogTitleTranslationMap = new Map Date: Mon, 26 May 2025 09:51:42 +0300 Subject: [PATCH 07/10] Moved TbServiceInfoProvider interface to separate module --- application/pom.xml | 4 + common/coap-server/pom.xml | 4 + common/edqs/pom.xml | 4 + common/pom.xml | 1 + common/queue/pom.xml | 4 + common/service-info-api/pom.xml | 73 +++++++++++++++++++ .../discovery/TbServiceInfoProvider.java | 0 common/transport/transport-api/pom.xml | 4 + dao/pom.xml | 2 +- pom.xml | 5 ++ 10 files changed, 100 insertions(+), 1 deletion(-) create mode 100644 common/service-info-api/pom.xml rename common/{queue => service-info-api}/src/main/java/org/thingsboard/server/queue/discovery/TbServiceInfoProvider.java (100%) diff --git a/application/pom.xml b/application/pom.xml index 557179e918..4f1d75cde2 100644 --- a/application/pom.xml +++ b/application/pom.xml @@ -128,6 +128,10 @@ org.thingsboard.common edqs
+ + org.thingsboard.common + service-info-api + org.thingsboard dao diff --git a/common/coap-server/pom.xml b/common/coap-server/pom.xml index fe6ae3c64c..44f16f7329 100644 --- a/common/coap-server/pom.xml +++ b/common/coap-server/pom.xml @@ -46,6 +46,10 @@ org.thingsboard.common data + + org.thingsboard.common + service-info-api + org.thingsboard.common.transport transport-api diff --git a/common/edqs/pom.xml b/common/edqs/pom.xml index 09181e796e..3b72a26442 100644 --- a/common/edqs/pom.xml +++ b/common/edqs/pom.xml @@ -68,6 +68,10 @@ org.thingsboard.common queue + + org.thingsboard.common + service-info-api + org.springframework.boot spring-boot-starter-web diff --git a/common/pom.xml b/common/pom.xml index 6afae4e378..514e3eea38 100644 --- a/common/pom.xml +++ b/common/pom.xml @@ -50,6 +50,7 @@ version-control script edqs + service-info-api diff --git a/common/queue/pom.xml b/common/queue/pom.xml index 2e4b14e282..e689cb3dc5 100644 --- a/common/queue/pom.xml +++ b/common/queue/pom.xml @@ -60,6 +60,10 @@ org.thingsboard.common cluster-api + + org.thingsboard.common + service-info-api + org.apache.kafka kafka-clients diff --git a/common/service-info-api/pom.xml b/common/service-info-api/pom.xml new file mode 100644 index 0000000000..c1eff7c468 --- /dev/null +++ b/common/service-info-api/pom.xml @@ -0,0 +1,73 @@ + + + 4.0.0 + + org.thingsboard + 4.1.0-SNAPSHOT + common + + org.thingsboard.common + service-info-api + jar + + Thingsboard Server Service Info Provider API + https://thingsboard.io + + + UTF-8 + ${basedir}/../.. + + + + + org.thingsboard.common + proto + + + org.thingsboard.common + message + + + + + + + org.apache.maven.plugins + maven-source-plugin + + + attach-sources + + jar + + + + + + org.apache.maven.plugins + maven-deploy-plugin + + false + + + + + + diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/TbServiceInfoProvider.java b/common/service-info-api/src/main/java/org/thingsboard/server/queue/discovery/TbServiceInfoProvider.java similarity index 100% rename from common/queue/src/main/java/org/thingsboard/server/queue/discovery/TbServiceInfoProvider.java rename to common/service-info-api/src/main/java/org/thingsboard/server/queue/discovery/TbServiceInfoProvider.java diff --git a/common/transport/transport-api/pom.xml b/common/transport/transport-api/pom.xml index 2cd41064da..064c5f8ea2 100644 --- a/common/transport/transport-api/pom.xml +++ b/common/transport/transport-api/pom.xml @@ -60,6 +60,10 @@ org.thingsboard.common util + + org.thingsboard.common + service-info-api + com.google.code.gson gson diff --git a/dao/pom.xml b/dao/pom.xml index 67b5a36c92..a11e4e4fe3 100644 --- a/dao/pom.xml +++ b/dao/pom.xml @@ -61,7 +61,7 @@ org.thingsboard.common - queue + service-info-api com.networknt diff --git a/pom.xml b/pom.xml index c57239c965..f8137196e4 100755 --- a/pom.xml +++ b/pom.xml @@ -964,6 +964,11 @@ cluster-api ${project.version} + + org.thingsboard.common + service-info-api + ${project.version} + org.thingsboard.rule-engine rule-engine-api From 92e81ad233f78c91dc03b5e244ef1ce7536cde55 Mon Sep 17 00:00:00 2001 From: dshvaika Date: Tue, 27 May 2025 16:42:21 +0300 Subject: [PATCH 08/10] renamed module service-info-api to discovery-api && fixed LimitedApi typos --- application/pom.xml | 2 +- common/coap-server/pom.xml | 2 +- .../thingsboard/server/common/data/limit/LimitedApi.java | 8 ++++---- common/{service-info-api => discovery-api}/pom.xml | 2 +- .../server/queue/discovery/TbServiceInfoProvider.java | 0 common/edqs/pom.xml | 2 +- common/pom.xml | 2 +- common/queue/pom.xml | 2 +- common/transport/transport-api/pom.xml | 2 +- dao/pom.xml | 2 +- pom.xml | 2 +- 11 files changed, 13 insertions(+), 13 deletions(-) rename common/{service-info-api => discovery-api}/pom.xml (98%) rename common/{service-info-api => discovery-api}/src/main/java/org/thingsboard/server/queue/discovery/TbServiceInfoProvider.java (100%) diff --git a/application/pom.xml b/application/pom.xml index 4f1d75cde2..34218956fc 100644 --- a/application/pom.xml +++ b/application/pom.xml @@ -130,7 +130,7 @@ org.thingsboard.common - service-info-api + discovery-api org.thingsboard diff --git a/common/coap-server/pom.xml b/common/coap-server/pom.xml index 44f16f7329..b5ca9c1b9c 100644 --- a/common/coap-server/pom.xml +++ b/common/coap-server/pom.xml @@ -48,7 +48,7 @@ org.thingsboard.common - service-info-api + discovery-api org.thingsboard.common.transport diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApi.java b/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApi.java index 19eb4b510f..2a9991a1cc 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApi.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApi.java @@ -30,10 +30,10 @@ public enum LimitedApi { REST_REQUESTS_PER_TENANT(DefaultTenantProfileConfiguration::getTenantServerRestLimitsConfiguration, "REST API requests", true), REST_REQUESTS_PER_CUSTOMER(DefaultTenantProfileConfiguration::getCustomerServerRestLimitsConfiguration, "REST API requests per customer", false), WS_UPDATES_PER_SESSION(DefaultTenantProfileConfiguration::getWsUpdatesPerSessionRateLimit, "WS updates per session", true), - CASSANDRA_WRITE_QUERIES_CORE(DefaultTenantProfileConfiguration::getCassandraReadQueryTenantCoreRateLimits, "Rest API Cassandra write queries", true), - CASSANDRA_READ_QUERIES_CORE(DefaultTenantProfileConfiguration::getCassandraWriteQueryTenantCoreRateLimits, "Rest API and WS telemetry Cassandra read queries", true), - CASSANDRA_WRITE_QUERIES_RULE_ENGINE(DefaultTenantProfileConfiguration::getCassandraReadQueryTenantRuleEngineRateLimits, "Rule Engine telemetry Cassandra write queries", true), - CASSANDRA_READ_QUERIES_RULE_ENGINE(DefaultTenantProfileConfiguration::getCassandraWriteQueryTenantRuleEngineRateLimits, "Rule Engine telemetry Cassandra read queries", true), + CASSANDRA_WRITE_QUERIES_CORE(DefaultTenantProfileConfiguration::getCassandraWriteQueryTenantCoreRateLimits, "Rest API Cassandra write queries", true), + CASSANDRA_READ_QUERIES_CORE(DefaultTenantProfileConfiguration::getCassandraReadQueryTenantCoreRateLimits, "Rest API and WS telemetry Cassandra read queries", true), + CASSANDRA_WRITE_QUERIES_RULE_ENGINE(DefaultTenantProfileConfiguration::getCassandraWriteQueryTenantRuleEngineRateLimits, "Rule Engine telemetry Cassandra write queries", true), + CASSANDRA_READ_QUERIES_RULE_ENGINE(DefaultTenantProfileConfiguration::getCassandraReadQueryTenantRuleEngineRateLimits, "Rule Engine telemetry Cassandra read queries", true), CASSANDRA_READ_QUERIES_MONOLITH( LimitedApiUtil.merge( DefaultTenantProfileConfiguration::getCassandraReadQueryTenantCoreRateLimits, diff --git a/common/service-info-api/pom.xml b/common/discovery-api/pom.xml similarity index 98% rename from common/service-info-api/pom.xml rename to common/discovery-api/pom.xml index c1eff7c468..5852761c4a 100644 --- a/common/service-info-api/pom.xml +++ b/common/discovery-api/pom.xml @@ -24,7 +24,7 @@ common org.thingsboard.common - service-info-api + discovery-api jar Thingsboard Server Service Info Provider API diff --git a/common/service-info-api/src/main/java/org/thingsboard/server/queue/discovery/TbServiceInfoProvider.java b/common/discovery-api/src/main/java/org/thingsboard/server/queue/discovery/TbServiceInfoProvider.java similarity index 100% rename from common/service-info-api/src/main/java/org/thingsboard/server/queue/discovery/TbServiceInfoProvider.java rename to common/discovery-api/src/main/java/org/thingsboard/server/queue/discovery/TbServiceInfoProvider.java diff --git a/common/edqs/pom.xml b/common/edqs/pom.xml index 3b72a26442..e349ae3087 100644 --- a/common/edqs/pom.xml +++ b/common/edqs/pom.xml @@ -70,7 +70,7 @@ org.thingsboard.common - service-info-api + discovery-api org.springframework.boot diff --git a/common/pom.xml b/common/pom.xml index 514e3eea38..0915f47f1a 100644 --- a/common/pom.xml +++ b/common/pom.xml @@ -50,7 +50,7 @@ version-control script edqs - service-info-api + discovery-api diff --git a/common/queue/pom.xml b/common/queue/pom.xml index e689cb3dc5..f2f3091999 100644 --- a/common/queue/pom.xml +++ b/common/queue/pom.xml @@ -62,7 +62,7 @@ org.thingsboard.common - service-info-api + discovery-api org.apache.kafka diff --git a/common/transport/transport-api/pom.xml b/common/transport/transport-api/pom.xml index 064c5f8ea2..7f54dcbc85 100644 --- a/common/transport/transport-api/pom.xml +++ b/common/transport/transport-api/pom.xml @@ -62,7 +62,7 @@ org.thingsboard.common - service-info-api + discovery-api com.google.code.gson diff --git a/dao/pom.xml b/dao/pom.xml index a11e4e4fe3..9f411f8302 100644 --- a/dao/pom.xml +++ b/dao/pom.xml @@ -61,7 +61,7 @@ org.thingsboard.common - service-info-api + discovery-api com.networknt diff --git a/pom.xml b/pom.xml index f8137196e4..334be1850e 100755 --- a/pom.xml +++ b/pom.xml @@ -966,7 +966,7 @@ org.thingsboard.common - service-info-api + discovery-api ${project.version} From 45dee234c8e62f25cc43da6bb0794a7bb5b76570 Mon Sep 17 00:00:00 2001 From: dshvaika Date: Wed, 28 May 2025 12:31:04 +0300 Subject: [PATCH 09/10] fixes after review + new tests added --- .../TenantProfileControllerTest.java | 38 +++++++ .../server/common/data/limit/LimitedApi.java | 8 +- ...mitedApiEntry.java => RateLimitEntry.java} | 6 +- ...LimitedApiUtil.java => RateLimitUtil.java} | 18 ++-- .../DefaultTenantProfileConfiguration.java | 56 +++++----- .../common/data/limit/LimitedApiTest.java | 102 ++++++++++++++++++ ...piUtilTest.java => RateLimitUtilTest.java} | 16 +-- .../server/common/msg/tools/TbRateLimits.java | 8 +- .../CassandraBufferedRateReadExecutor.java | 9 +- .../CassandraBufferedRateWriteExecutor.java | 9 +- .../dao/service/RateLimitValidator.java | 4 +- .../util/AbstractBufferedRateExecutor.java | 57 +++++----- 12 files changed, 227 insertions(+), 104 deletions(-) rename common/data/src/main/java/org/thingsboard/server/common/data/limit/{LimitedApiEntry.java => RateLimitEntry.java} (79%) rename common/data/src/main/java/org/thingsboard/server/common/data/limit/{LimitedApiUtil.java => RateLimitUtil.java} (84%) create mode 100644 common/data/src/test/java/org/thingsboard/server/common/data/limit/LimitedApiTest.java rename common/data/src/test/java/org/thingsboard/server/common/data/limit/{LimitedApiUtilTest.java => RateLimitUtilTest.java} (86%) diff --git a/application/src/test/java/org/thingsboard/server/controller/TenantProfileControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/TenantProfileControllerTest.java index 0daa56728b..8f83f2c186 100644 --- a/application/src/test/java/org/thingsboard/server/controller/TenantProfileControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/TenantProfileControllerTest.java @@ -36,9 +36,11 @@ import org.thingsboard.server.common.data.queue.SubmitStrategyType; import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration; import org.thingsboard.server.common.data.tenant.profile.TenantProfileData; import org.thingsboard.server.common.data.tenant.profile.TenantProfileQueueConfiguration; +import org.thingsboard.server.common.data.validation.RateLimit; import org.thingsboard.server.dao.service.DaoSqlTest; import org.thingsboard.server.queue.TbQueueCallback; +import java.lang.reflect.Field; import java.util.ArrayList; import java.util.Collections; import java.util.List; @@ -314,6 +316,42 @@ public class TenantProfileControllerTest extends AbstractControllerTest { Assert.assertEquals(1, pageData.getTotalElements()); } + @Test + public void testRateLimitValidationAllFields() throws Exception { + loginSysAdmin(); + Mockito.reset(tbClusterService); + + List failedFields = new ArrayList<>(); + + for (Field field : DefaultTenantProfileConfiguration.class.getDeclaredFields()) { + RateLimit rateLimit = field.getAnnotation(RateLimit.class); + if (rateLimit == null) continue; + + String fieldName = field.getName(); + String expectedLabel = rateLimit.fieldName(); + + + TenantProfile tenantProfile = createTenantProfile("Invalid RateLimit - " + fieldName); + DefaultTenantProfileConfiguration config = (DefaultTenantProfileConfiguration) tenantProfile.getProfileData().getConfiguration(); + + field.setAccessible(true); + field.set(config, "10:1,10:1"); // Set invalid duplicate value + + try { + doPost("/api/tenantProfile", tenantProfile) + .andExpect(status().isBadRequest()) + .andExpect(statusReason(containsString(expectedLabel + " rate limit has duplicate 'Per seconds' configuration."))); + } catch (AssertionError e) { + failedFields.add(fieldName + " (label: " + expectedLabel + ")"); + } + } + + if (!failedFields.isEmpty()) { + throw new AssertionError("RateLimit validation failed for fields: " + String.join(", ", failedFields)); + } + testBroadcastEntityStateChangeEventNeverTenantProfile(); + } + private TenantProfile createTenantProfile(String name) { TenantProfile tenantProfile = new TenantProfile(); tenantProfile.setName(name); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApi.java b/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApi.java index 2a9991a1cc..a2582db9f3 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApi.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApi.java @@ -21,6 +21,7 @@ import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileCon import java.util.Optional; import java.util.function.Function; +@Getter public enum LimitedApi { ENTITY_EXPORT(DefaultTenantProfileConfiguration::getTenantEntityExportRateLimit, "entity version creation", true), @@ -35,11 +36,11 @@ public enum LimitedApi { CASSANDRA_WRITE_QUERIES_RULE_ENGINE(DefaultTenantProfileConfiguration::getCassandraWriteQueryTenantRuleEngineRateLimits, "Rule Engine telemetry Cassandra write queries", true), CASSANDRA_READ_QUERIES_RULE_ENGINE(DefaultTenantProfileConfiguration::getCassandraReadQueryTenantRuleEngineRateLimits, "Rule Engine telemetry Cassandra read queries", true), CASSANDRA_READ_QUERIES_MONOLITH( - LimitedApiUtil.merge( + RateLimitUtil.merge( DefaultTenantProfileConfiguration::getCassandraReadQueryTenantCoreRateLimits, DefaultTenantProfileConfiguration::getCassandraReadQueryTenantRuleEngineRateLimits), "Telemetry read queries", true), CASSANDRA_WRITE_QUERIES_MONOLITH( - LimitedApiUtil.merge( + RateLimitUtil.merge( DefaultTenantProfileConfiguration::getCassandraWriteQueryTenantCoreRateLimits, DefaultTenantProfileConfiguration::getCassandraWriteQueryTenantRuleEngineRateLimits), "Telemetry write queries", true), EDGE_EVENTS(DefaultTenantProfileConfiguration::getEdgeEventRateLimits, "Edge events", true), @@ -58,11 +59,8 @@ public enum LimitedApi { CALCULATED_FIELD_DEBUG_EVENTS("calculated field debug events", true); private final Function configExtractor; - @Getter private final boolean perTenant; - @Getter private final boolean refillRateLimitIntervally; - @Getter private final String label; LimitedApi(Function configExtractor, String label, boolean perTenant) { diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApiEntry.java b/common/data/src/main/java/org/thingsboard/server/common/data/limit/RateLimitEntry.java similarity index 79% rename from common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApiEntry.java rename to common/data/src/main/java/org/thingsboard/server/common/data/limit/RateLimitEntry.java index 16084b7af9..566ff96c75 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApiEntry.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/limit/RateLimitEntry.java @@ -15,11 +15,11 @@ */ package org.thingsboard.server.common.data.limit; -public record LimitedApiEntry(long capacity, long durationSeconds) { +public record RateLimitEntry(long capacity, long durationSeconds) { - public static LimitedApiEntry parse(String s) { + public static RateLimitEntry parse(String s) { String[] parts = s.split(":"); - return new LimitedApiEntry(Long.parseLong(parts[0]), Long.parseLong(parts[1])); + return new RateLimitEntry(Long.parseLong(parts[0]), Long.parseLong(parts[1])); } @Override diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApiUtil.java b/common/data/src/main/java/org/thingsboard/server/common/data/limit/RateLimitUtil.java similarity index 84% rename from common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApiUtil.java rename to common/data/src/main/java/org/thingsboard/server/common/data/limit/RateLimitUtil.java index 9c99d065d4..f2ee164dcd 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApiUtil.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/limit/RateLimitUtil.java @@ -28,14 +28,14 @@ import java.util.Set; import java.util.function.Function; import java.util.stream.Collectors; -public class LimitedApiUtil { +public class RateLimitUtil { - public static List parseConfig(String config) { + public static List parseConfig(String config) { if (config == null || config.isEmpty()) { return Collections.emptyList(); } return Arrays.stream(config.split(",")) - .map(LimitedApiEntry::parse) + .map(RateLimitEntry::parse) .toList(); } @@ -45,18 +45,18 @@ public class LimitedApiUtil { return config -> { String config1 = configExtractor1.apply(config); String config2 = configExtractor2.apply(config); - return LimitedApiUtil.mergeStrConfigs(config1, config2); // merges the configs + return RateLimitUtil.mergeStrConfigs(config1, config2); // merges the configs }; } private static String mergeStrConfigs(String firstConfig, String secondConfig) { - List all = new ArrayList<>(); + List all = new ArrayList<>(); all.addAll(parseConfig(firstConfig)); all.addAll(parseConfig(secondConfig)); Map merged = new HashMap<>(); - for (LimitedApiEntry entry : all) { + for (RateLimitEntry entry : all) { merged.merge(entry.durationSeconds(), entry.capacity(), Long::sum); } @@ -67,9 +67,9 @@ public class LimitedApiUtil { } public static boolean isValid(String configStr) { - List limitedApiEntries = parseConfig(configStr); + List limitedApiEntries = parseConfig(configStr); Set distinctDurations = new HashSet<>(); - for (LimitedApiEntry entry : limitedApiEntries) { + for (RateLimitEntry entry : limitedApiEntries) { if (!distinctDurations.add(entry.durationSeconds())) { return false; } @@ -85,7 +85,7 @@ public class LimitedApiUtil { Set distinctDurations = new HashSet<>(); return parseConfig(configStr).stream() .filter(entry -> distinctDurations.add(entry.durationSeconds())) - .map(LimitedApiEntry::toString) + .map(RateLimitEntry::toString) .collect(Collectors.joining(",")); } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java index 00fecfe800..a4ff47c340 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java @@ -24,7 +24,7 @@ import lombok.NoArgsConstructor; import org.thingsboard.server.common.data.ApiUsageRecordKey; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.TenantProfileType; -import org.thingsboard.server.common.data.limit.LimitedApiUtil; +import org.thingsboard.server.common.data.limit.RateLimitUtil; import org.thingsboard.server.common.data.validation.RateLimit; import java.io.Serial; @@ -238,41 +238,41 @@ public class DefaultTenantProfileConfiguration implements TenantProfileConfigura @Deprecated(forRemoval = true, since = "4.1") public void deduplicateRateLimitsConfigs() { - this.transportTenantMsgRateLimit = LimitedApiUtil.deduplicateByDuration(transportTenantMsgRateLimit); - this.transportTenantTelemetryMsgRateLimit = LimitedApiUtil.deduplicateByDuration(transportTenantTelemetryMsgRateLimit); - this.transportTenantTelemetryDataPointsRateLimit = LimitedApiUtil.deduplicateByDuration(transportTenantTelemetryDataPointsRateLimit); + this.transportTenantMsgRateLimit = RateLimitUtil.deduplicateByDuration(transportTenantMsgRateLimit); + this.transportTenantTelemetryMsgRateLimit = RateLimitUtil.deduplicateByDuration(transportTenantTelemetryMsgRateLimit); + this.transportTenantTelemetryDataPointsRateLimit = RateLimitUtil.deduplicateByDuration(transportTenantTelemetryDataPointsRateLimit); - this.transportDeviceMsgRateLimit = LimitedApiUtil.deduplicateByDuration(transportDeviceMsgRateLimit); - this.transportDeviceTelemetryMsgRateLimit = LimitedApiUtil.deduplicateByDuration(transportDeviceTelemetryMsgRateLimit); - this.transportDeviceTelemetryDataPointsRateLimit = LimitedApiUtil.deduplicateByDuration(transportDeviceTelemetryDataPointsRateLimit); + this.transportDeviceMsgRateLimit = RateLimitUtil.deduplicateByDuration(transportDeviceMsgRateLimit); + this.transportDeviceTelemetryMsgRateLimit = RateLimitUtil.deduplicateByDuration(transportDeviceTelemetryMsgRateLimit); + this.transportDeviceTelemetryDataPointsRateLimit = RateLimitUtil.deduplicateByDuration(transportDeviceTelemetryDataPointsRateLimit); - this.transportGatewayMsgRateLimit = LimitedApiUtil.deduplicateByDuration(transportGatewayMsgRateLimit); - this.transportGatewayTelemetryMsgRateLimit = LimitedApiUtil.deduplicateByDuration(transportGatewayTelemetryMsgRateLimit); - this.transportGatewayTelemetryDataPointsRateLimit = LimitedApiUtil.deduplicateByDuration(transportGatewayTelemetryDataPointsRateLimit); + this.transportGatewayMsgRateLimit = RateLimitUtil.deduplicateByDuration(transportGatewayMsgRateLimit); + this.transportGatewayTelemetryMsgRateLimit = RateLimitUtil.deduplicateByDuration(transportGatewayTelemetryMsgRateLimit); + this.transportGatewayTelemetryDataPointsRateLimit = RateLimitUtil.deduplicateByDuration(transportGatewayTelemetryDataPointsRateLimit); - this.transportGatewayDeviceMsgRateLimit = LimitedApiUtil.deduplicateByDuration(transportGatewayDeviceMsgRateLimit); - this.transportGatewayDeviceTelemetryMsgRateLimit = LimitedApiUtil.deduplicateByDuration(transportGatewayDeviceTelemetryMsgRateLimit); - this.transportGatewayDeviceTelemetryDataPointsRateLimit = LimitedApiUtil.deduplicateByDuration(transportGatewayDeviceTelemetryDataPointsRateLimit); + this.transportGatewayDeviceMsgRateLimit = RateLimitUtil.deduplicateByDuration(transportGatewayDeviceMsgRateLimit); + this.transportGatewayDeviceTelemetryMsgRateLimit = RateLimitUtil.deduplicateByDuration(transportGatewayDeviceTelemetryMsgRateLimit); + this.transportGatewayDeviceTelemetryDataPointsRateLimit = RateLimitUtil.deduplicateByDuration(transportGatewayDeviceTelemetryDataPointsRateLimit); - this.tenantEntityExportRateLimit = LimitedApiUtil.deduplicateByDuration(tenantEntityExportRateLimit); - this.tenantEntityImportRateLimit = LimitedApiUtil.deduplicateByDuration(tenantEntityImportRateLimit); - this.tenantNotificationRequestsRateLimit = LimitedApiUtil.deduplicateByDuration(tenantNotificationRequestsRateLimit); - this.tenantNotificationRequestsPerRuleRateLimit = LimitedApiUtil.deduplicateByDuration(tenantNotificationRequestsPerRuleRateLimit); + this.tenantEntityExportRateLimit = RateLimitUtil.deduplicateByDuration(tenantEntityExportRateLimit); + this.tenantEntityImportRateLimit = RateLimitUtil.deduplicateByDuration(tenantEntityImportRateLimit); + this.tenantNotificationRequestsRateLimit = RateLimitUtil.deduplicateByDuration(tenantNotificationRequestsRateLimit); + this.tenantNotificationRequestsPerRuleRateLimit = RateLimitUtil.deduplicateByDuration(tenantNotificationRequestsPerRuleRateLimit); - this.cassandraReadQueryTenantCoreRateLimits = LimitedApiUtil.deduplicateByDuration(cassandraReadQueryTenantCoreRateLimits); - this.cassandraWriteQueryTenantCoreRateLimits = LimitedApiUtil.deduplicateByDuration(cassandraWriteQueryTenantCoreRateLimits); - this.cassandraReadQueryTenantRuleEngineRateLimits = LimitedApiUtil.deduplicateByDuration(cassandraReadQueryTenantRuleEngineRateLimits); - this.cassandraWriteQueryTenantRuleEngineRateLimits = LimitedApiUtil.deduplicateByDuration(cassandraWriteQueryTenantRuleEngineRateLimits); + this.cassandraReadQueryTenantCoreRateLimits = RateLimitUtil.deduplicateByDuration(cassandraReadQueryTenantCoreRateLimits); + this.cassandraWriteQueryTenantCoreRateLimits = RateLimitUtil.deduplicateByDuration(cassandraWriteQueryTenantCoreRateLimits); + this.cassandraReadQueryTenantRuleEngineRateLimits = RateLimitUtil.deduplicateByDuration(cassandraReadQueryTenantRuleEngineRateLimits); + this.cassandraWriteQueryTenantRuleEngineRateLimits = RateLimitUtil.deduplicateByDuration(cassandraWriteQueryTenantRuleEngineRateLimits); - this.edgeEventRateLimits = LimitedApiUtil.deduplicateByDuration(edgeEventRateLimits); - this.edgeEventRateLimitsPerEdge = LimitedApiUtil.deduplicateByDuration(edgeEventRateLimitsPerEdge); - this.edgeUplinkMessagesRateLimits = LimitedApiUtil.deduplicateByDuration(edgeUplinkMessagesRateLimits); - this.edgeUplinkMessagesRateLimitsPerEdge = LimitedApiUtil.deduplicateByDuration(edgeUplinkMessagesRateLimitsPerEdge); + this.edgeEventRateLimits = RateLimitUtil.deduplicateByDuration(edgeEventRateLimits); + this.edgeEventRateLimitsPerEdge = RateLimitUtil.deduplicateByDuration(edgeEventRateLimitsPerEdge); + this.edgeUplinkMessagesRateLimits = RateLimitUtil.deduplicateByDuration(edgeUplinkMessagesRateLimits); + this.edgeUplinkMessagesRateLimitsPerEdge = RateLimitUtil.deduplicateByDuration(edgeUplinkMessagesRateLimitsPerEdge); - this.wsUpdatesPerSessionRateLimit = LimitedApiUtil.deduplicateByDuration(wsUpdatesPerSessionRateLimit); + this.wsUpdatesPerSessionRateLimit = RateLimitUtil.deduplicateByDuration(wsUpdatesPerSessionRateLimit); - this.tenantServerRestLimitsConfiguration = LimitedApiUtil.deduplicateByDuration(tenantServerRestLimitsConfiguration); - this.customerServerRestLimitsConfiguration = LimitedApiUtil.deduplicateByDuration(customerServerRestLimitsConfiguration); + this.tenantServerRestLimitsConfiguration = RateLimitUtil.deduplicateByDuration(tenantServerRestLimitsConfiguration); + this.customerServerRestLimitsConfiguration = RateLimitUtil.deduplicateByDuration(customerServerRestLimitsConfiguration); } } diff --git a/common/data/src/test/java/org/thingsboard/server/common/data/limit/LimitedApiTest.java b/common/data/src/test/java/org/thingsboard/server/common/data/limit/LimitedApiTest.java new file mode 100644 index 0000000000..d650017fd4 --- /dev/null +++ b/common/data/src/test/java/org/thingsboard/server/common/data/limit/LimitedApiTest.java @@ -0,0 +1,102 @@ +/** + * Copyright © 2016-2025 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.data.limit; + +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration; + +import java.util.Arrays; +import java.util.Map; +import java.util.Set; +import java.util.stream.Collectors; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.Mockito.clearInvocations; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; + +class LimitedApiTest { + + private DefaultTenantProfileConfiguration config; + + @BeforeEach + void setUp() { + config = mock(DefaultTenantProfileConfiguration.class); + } + + @Test + void testCorrectConfigExtractorsUsed() { + Map verifierMap = Map.ofEntries( + Map.entry(LimitedApi.ENTITY_EXPORT, () -> + verify(config).getTenantEntityExportRateLimit()), + Map.entry(LimitedApi.ENTITY_IMPORT, () -> + verify(config).getTenantEntityImportRateLimit()), + Map.entry(LimitedApi.NOTIFICATION_REQUESTS, () -> + verify(config).getTenantNotificationRequestsRateLimit()), + Map.entry(LimitedApi.NOTIFICATION_REQUESTS_PER_RULE, () -> + verify(config).getTenantNotificationRequestsPerRuleRateLimit()), + Map.entry(LimitedApi.REST_REQUESTS_PER_TENANT, () -> + verify(config).getTenantServerRestLimitsConfiguration()), + Map.entry(LimitedApi.REST_REQUESTS_PER_CUSTOMER, () -> + verify(config).getCustomerServerRestLimitsConfiguration()), + Map.entry(LimitedApi.WS_UPDATES_PER_SESSION, () -> + verify(config).getWsUpdatesPerSessionRateLimit()), + Map.entry(LimitedApi.CASSANDRA_WRITE_QUERIES_CORE, () -> + verify(config).getCassandraWriteQueryTenantCoreRateLimits()), + Map.entry(LimitedApi.CASSANDRA_READ_QUERIES_CORE, () -> + verify(config).getCassandraReadQueryTenantCoreRateLimits()), + Map.entry(LimitedApi.CASSANDRA_WRITE_QUERIES_RULE_ENGINE, () -> + verify(config).getCassandraWriteQueryTenantRuleEngineRateLimits()), + Map.entry(LimitedApi.CASSANDRA_READ_QUERIES_RULE_ENGINE, () -> + verify(config).getCassandraReadQueryTenantRuleEngineRateLimits()), + Map.entry(LimitedApi.CASSANDRA_READ_QUERIES_MONOLITH, () -> { + verify(config).getCassandraReadQueryTenantCoreRateLimits(); + verify(config).getCassandraReadQueryTenantRuleEngineRateLimits(); + }), + Map.entry(LimitedApi.CASSANDRA_WRITE_QUERIES_MONOLITH, () -> { + verify(config).getCassandraWriteQueryTenantCoreRateLimits(); + verify(config).getCassandraWriteQueryTenantRuleEngineRateLimits(); + }), + Map.entry(LimitedApi.EDGE_EVENTS, () -> + verify(config).getEdgeEventRateLimits()), + Map.entry(LimitedApi.EDGE_EVENTS_PER_EDGE, () -> + verify(config).getEdgeEventRateLimitsPerEdge()), + Map.entry(LimitedApi.EDGE_UPLINK_MESSAGES, () -> + verify(config).getEdgeUplinkMessagesRateLimits()), + Map.entry(LimitedApi.EDGE_UPLINK_MESSAGES_PER_EDGE, () -> + verify(config).getEdgeUplinkMessagesRateLimitsPerEdge()) + ); + + Set expected = verifierMap.keySet(); + Set actual = Arrays.stream(LimitedApi.values()) + .filter(api -> api.getConfigExtractor() != null) + .collect(Collectors.toSet()); + + assertThat(expected) + .as("Verifier map should cover all LimitedApis with extractors") + .containsExactlyInAnyOrderElementsOf(actual); + + for (Map.Entry entry : verifierMap.entrySet()) { + LimitedApi api = entry.getKey(); + api.getLimitConfig(config); + entry.getValue().run(); + clearInvocations(config); + } + } + + +} \ No newline at end of file diff --git a/common/data/src/test/java/org/thingsboard/server/common/data/limit/LimitedApiUtilTest.java b/common/data/src/test/java/org/thingsboard/server/common/data/limit/RateLimitUtilTest.java similarity index 86% rename from common/data/src/test/java/org/thingsboard/server/common/data/limit/LimitedApiUtilTest.java rename to common/data/src/test/java/org/thingsboard/server/common/data/limit/RateLimitUtilTest.java index c275df790c..a410785126 100644 --- a/common/data/src/test/java/org/thingsboard/server/common/data/limit/LimitedApiUtilTest.java +++ b/common/data/src/test/java/org/thingsboard/server/common/data/limit/RateLimitUtilTest.java @@ -24,12 +24,12 @@ import java.util.function.Function; import static org.assertj.core.api.Assertions.assertThat; -class LimitedApiUtilTest { +class RateLimitUtilTest { @Test @DisplayName("LimitedApiUtil should parse single entry correctly") void testParseSingleEntry() { - List entries = LimitedApiUtil.parseConfig("100:60"); + List entries = RateLimitUtil.parseConfig("100:60"); assertThat(entries).hasSize(1); assertThat(entries.get(0).capacity()).isEqualTo(100); @@ -39,7 +39,7 @@ class LimitedApiUtilTest { @Test @DisplayName("LimitedApiUtil should parse multiple entries correctly") void testParseMultipleEntries() { - List entries = LimitedApiUtil.parseConfig("100:60,200:30"); + List entries = RateLimitUtil.parseConfig("100:60,200:30"); assertThat(entries).hasSize(2); assertThat(entries.get(0).capacity()).isEqualTo(100); @@ -51,8 +51,8 @@ class LimitedApiUtilTest { @Test @DisplayName("LimitedApiUtil should return empty list for null or empty config") void testParseEmptyConfig() { - assertThat(LimitedApiUtil.parseConfig(null)).isEmpty(); - assertThat(LimitedApiUtil.parseConfig("")).isEmpty(); + assertThat(RateLimitUtil.parseConfig(null)).isEmpty(); + assertThat(RateLimitUtil.parseConfig("")).isEmpty(); } @Test @@ -64,7 +64,7 @@ class LimitedApiUtilTest { // Fake config instance (not used directly in lambda logic) DefaultTenantProfileConfiguration config = new DefaultTenantProfileConfiguration(); - String result = LimitedApiUtil.merge(extractor1, extractor2).apply(config); + String result = RateLimitUtil.merge(extractor1, extractor2).apply(config); // Should be: 300:60 (100+200), 50:30, 25:10 assertThat(result).isEqualTo("25:10,50:30,300:60"); @@ -78,7 +78,7 @@ class LimitedApiUtilTest { // Fake config instance (not used directly in lambda logic) DefaultTenantProfileConfiguration config = new DefaultTenantProfileConfiguration(); - String result = LimitedApiUtil.merge(extractor1, extractor2).apply(config); + String result = RateLimitUtil.merge(extractor1, extractor2).apply(config); assertThat(result).isEqualTo("100:60"); } @@ -91,7 +91,7 @@ class LimitedApiUtilTest { // Fake config instance (not used directly in lambda logic) DefaultTenantProfileConfiguration config = new DefaultTenantProfileConfiguration(); - String result = LimitedApiUtil.merge(extractor1, extractor2).apply(config); + String result = RateLimitUtil.merge(extractor1, extractor2).apply(config); assertThat(result).isEqualTo("200:10,100:60"); } diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/tools/TbRateLimits.java b/common/message/src/main/java/org/thingsboard/server/common/msg/tools/TbRateLimits.java index c65c90706a..182cb21ab8 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/tools/TbRateLimits.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/tools/TbRateLimits.java @@ -21,8 +21,8 @@ import io.github.bucket4j.Bucket; import io.github.bucket4j.local.LocalBucket; import io.github.bucket4j.local.LocalBucketBuilder; import lombok.Getter; -import org.thingsboard.server.common.data.limit.LimitedApiEntry; -import org.thingsboard.server.common.data.limit.LimitedApiUtil; +import org.thingsboard.server.common.data.limit.RateLimitEntry; +import org.thingsboard.server.common.data.limit.RateLimitUtil; import java.time.Duration; import java.util.List; @@ -41,12 +41,12 @@ public class TbRateLimits { } public TbRateLimits(String limitsConfiguration, boolean refillIntervally) { - List limitedApiEntries = LimitedApiUtil.parseConfig(limitsConfiguration); + List limitedApiEntries = RateLimitUtil.parseConfig(limitsConfiguration); if (limitedApiEntries.isEmpty()) { throw new IllegalArgumentException("Failed to parse rate limits configuration: " + limitsConfiguration); } LocalBucketBuilder localBucket = Bucket.builder(); - for (LimitedApiEntry entry : limitedApiEntries) { + for (RateLimitEntry entry : limitedApiEntries) { BandwidthBuilder.BandwidthBuilderRefillStage bandwidthBuilder = Bandwidth.builder().capacity(entry.capacity()); Bandwidth bandwidth = refillIntervally ? bandwidthBuilder.refillIntervally(entry.capacity(), Duration.ofSeconds(entry.durationSeconds())).build() : diff --git a/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraBufferedRateReadExecutor.java b/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraBufferedRateReadExecutor.java index 98cfef7f7d..8698429b60 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraBufferedRateReadExecutor.java +++ b/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraBufferedRateReadExecutor.java @@ -53,8 +53,8 @@ public class CassandraBufferedRateReadExecutor extends AbstractBufferedRateExecu @Autowired EntityService entityService, @Autowired RateLimitService rateLimitService, @Autowired(required = false) TbServiceInfoProvider serviceInfoProvider) { - super(queueLimit, concurrencyLimit, maxWaitTime, dispatcherThreads, callbackThreads, pollMs, printQueriesFreq, statsFactory, - entityService, rateLimitService, serviceInfoProvider, printTenantNames); + super(queueLimit, concurrencyLimit, maxWaitTime, dispatcherThreads, callbackThreads, pollMs, printQueriesFreq, + BufferedRateExecutorType.READ, entityService, rateLimitService, serviceInfoProvider, statsFactory, printTenantNames); } @Scheduled(fixedDelayString = "${cassandra.query.rate_limit_print_interval_ms}") @@ -68,11 +68,6 @@ public class CassandraBufferedRateReadExecutor extends AbstractBufferedRateExecu super.stop(); } - @Override - protected BufferedRateExecutorType getBufferedRateExecutorType() { - return BufferedRateExecutorType.READ; - } - @Override protected SettableFuture create() { return SettableFuture.create(); diff --git a/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraBufferedRateWriteExecutor.java b/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraBufferedRateWriteExecutor.java index 1fb4b8dea5..6eb313be96 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraBufferedRateWriteExecutor.java +++ b/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraBufferedRateWriteExecutor.java @@ -53,8 +53,8 @@ public class CassandraBufferedRateWriteExecutor extends AbstractBufferedRateExec @Autowired EntityService entityService, @Autowired RateLimitService rateLimitService, @Autowired(required = false) TbServiceInfoProvider serviceInfoProvider) { - super(queueLimit, concurrencyLimit, maxWaitTime, dispatcherThreads, callbackThreads, pollMs, printQueriesFreq, statsFactory, - entityService, rateLimitService, serviceInfoProvider, printTenantNames); + super(queueLimit, concurrencyLimit, maxWaitTime, dispatcherThreads, callbackThreads, pollMs, printQueriesFreq, + BufferedRateExecutorType.WRITE, entityService, rateLimitService, serviceInfoProvider, statsFactory, printTenantNames); } @Scheduled(fixedDelayString = "${cassandra.query.rate_limit_print_interval_ms}") @@ -68,11 +68,6 @@ public class CassandraBufferedRateWriteExecutor extends AbstractBufferedRateExec super.stop(); } - @Override - protected BufferedRateExecutorType getBufferedRateExecutorType() { - return BufferedRateExecutorType.WRITE; - } - @Override protected SettableFuture create() { return SettableFuture.create(); diff --git a/dao/src/main/java/org/thingsboard/server/dao/service/RateLimitValidator.java b/dao/src/main/java/org/thingsboard/server/dao/service/RateLimitValidator.java index 0703f8b65d..bf776d4dc7 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/service/RateLimitValidator.java +++ b/dao/src/main/java/org/thingsboard/server/dao/service/RateLimitValidator.java @@ -18,7 +18,7 @@ package org.thingsboard.server.dao.service; import jakarta.validation.ConstraintValidator; import jakarta.validation.ConstraintValidatorContext; import lombok.extern.slf4j.Slf4j; -import org.thingsboard.server.common.data.limit.LimitedApiUtil; +import org.thingsboard.server.common.data.limit.RateLimitUtil; import org.thingsboard.server.common.data.validation.RateLimit; @Slf4j @@ -26,7 +26,7 @@ public class RateLimitValidator implements ConstraintValidator> queue; private final ExecutorService dispatcherExecutor; private final ExecutorService callbackExecutor; @@ -80,29 +81,32 @@ public abstract class AbstractBufferedRateExecutor tenantNamesCache = new HashMap<>(); + private final LimitedApi myLimitedApi; + public AbstractBufferedRateExecutor(int queueLimit, int concurrencyLimit, long maxWaitTime, int dispatcherThreads, - int callbackThreads, long pollMs, int printQueriesFreq, StatsFactory statsFactory, - EntityService entityService, RateLimitService rateLimitService, TbServiceInfoProvider serviceInfoProvider, boolean printTenantNames) { + int callbackThreads, long pollMs, int printQueriesFreq, BufferedRateExecutorType executorType, + EntityService entityService, RateLimitService rateLimitService, TbServiceInfoProvider serviceInfoProvider, + StatsFactory statsFactory, boolean printTenantNames) { this.maxWaitTime = maxWaitTime; this.pollMs = pollMs; + this.bufferName = executorType.getDisplayName(); this.concurrencyLimit = concurrencyLimit; this.printQueriesFreq = printQueriesFreq; this.queue = new LinkedBlockingDeque<>(queueLimit); - this.dispatcherExecutor = Executors.newFixedThreadPool(dispatcherThreads, ThingsBoardThreadFactory.forName("nosql-" + getBufferName() + "-dispatcher")); - this.callbackExecutor = ThingsBoardExecutors.newWorkStealingPool(callbackThreads, "nosql-" + getBufferName() + "-callback"); - this.timeoutExecutor = ThingsBoardExecutors.newSingleThreadScheduledExecutor("nosql-" + getBufferName() + "-timeout"); + this.dispatcherExecutor = Executors.newFixedThreadPool(dispatcherThreads, ThingsBoardThreadFactory.forName("nosql-" + bufferName + "-dispatcher")); + this.callbackExecutor = ThingsBoardExecutors.newWorkStealingPool(callbackThreads, "nosql-" + bufferName + "-callback"); + this.timeoutExecutor = ThingsBoardExecutors.newSingleThreadScheduledExecutor("nosql-" + bufferName + "-timeout"); this.stats = new BufferedRateExecutorStats(statsFactory); - String concurrencyLevelKey = StatsType.RATE_EXECUTOR.getName() + "." + CONCURRENCY_LEVEL + getBufferName(); //metric name may change with buffer name suffix + String concurrencyLevelKey = StatsType.RATE_EXECUTOR.getName() + "." + CONCURRENCY_LEVEL + bufferName; //metric name may change with buffer name suffix this.concurrencyLevel = statsFactory.createGauge(concurrencyLevelKey, new AtomicInteger(0)); this.entityService = entityService; this.rateLimitService = rateLimitService; - this.serviceInfoProvider = serviceInfoProvider; + this.myLimitedApi = resolveLimitedApi(serviceInfoProvider, executorType); this.printTenantNames = printTenantNames; for (int i = 0; i < dispatcherThreads; i++) { @@ -118,14 +122,14 @@ public abstract class AbstractBufferedRateExecutor execute(AsyncTaskContext taskCtx); - private String getBufferName() { - return getBufferedRateExecutorType().getDisplayName(); - } - - protected abstract BufferedRateExecutorType getBufferedRateExecutorType(); - private void dispatch() { - log.info("[{}] Buffered rate executor thread started", getBufferName()); + log.info("[{}] Buffered rate executor thread started", bufferName); while (!Thread.interrupted()) { int curLvl = concurrencyLevel.get(); AsyncTaskContext taskCtx = null; @@ -190,7 +185,7 @@ public abstract class AbstractBufferedRateExecutor= printQueriesFreq) { printQueriesIdx.set(0); String query = queryToString(finalTaskCtx); - log.info("[{}][{}] Cassandra query: {}", getBufferName(), taskCtx.getId(), query); + log.info("[{}][{}] Cassandra query: {}", bufferName, taskCtx.getId(), query); } } logTask("Processing", finalTaskCtx); @@ -243,7 +238,7 @@ public abstract class AbstractBufferedRateExecutor taskCtx) { @@ -319,7 +314,7 @@ public abstract class AbstractBufferedRateExecutor Date: Wed, 28 May 2025 12:38:13 +0300 Subject: [PATCH 10/10] fix typos --- .../server/common/data/limit/LimitedApiTest.java | 3 +-- .../dao/nosql/CassandraBufferedRateReadExecutor.java | 2 +- .../dao/nosql/CassandraBufferedRateWriteExecutor.java | 2 +- .../server/dao/util/AbstractBufferedRateExecutor.java | 7 +++---- 4 files changed, 6 insertions(+), 8 deletions(-) diff --git a/common/data/src/test/java/org/thingsboard/server/common/data/limit/LimitedApiTest.java b/common/data/src/test/java/org/thingsboard/server/common/data/limit/LimitedApiTest.java index d650017fd4..6c5e62e456 100644 --- a/common/data/src/test/java/org/thingsboard/server/common/data/limit/LimitedApiTest.java +++ b/common/data/src/test/java/org/thingsboard/server/common/data/limit/LimitedApiTest.java @@ -98,5 +98,4 @@ class LimitedApiTest { } } - -} \ No newline at end of file +} diff --git a/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraBufferedRateReadExecutor.java b/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraBufferedRateReadExecutor.java index 8698429b60..7befb1db62 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraBufferedRateReadExecutor.java +++ b/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraBufferedRateReadExecutor.java @@ -54,7 +54,7 @@ public class CassandraBufferedRateReadExecutor extends AbstractBufferedRateExecu @Autowired RateLimitService rateLimitService, @Autowired(required = false) TbServiceInfoProvider serviceInfoProvider) { super(queueLimit, concurrencyLimit, maxWaitTime, dispatcherThreads, callbackThreads, pollMs, printQueriesFreq, - BufferedRateExecutorType.READ, entityService, rateLimitService, serviceInfoProvider, statsFactory, printTenantNames); + BufferedRateExecutorType.READ, serviceInfoProvider, rateLimitService, statsFactory, entityService, printTenantNames); } @Scheduled(fixedDelayString = "${cassandra.query.rate_limit_print_interval_ms}") diff --git a/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraBufferedRateWriteExecutor.java b/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraBufferedRateWriteExecutor.java index 6eb313be96..9b7db4fae4 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraBufferedRateWriteExecutor.java +++ b/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraBufferedRateWriteExecutor.java @@ -54,7 +54,7 @@ public class CassandraBufferedRateWriteExecutor extends AbstractBufferedRateExec @Autowired RateLimitService rateLimitService, @Autowired(required = false) TbServiceInfoProvider serviceInfoProvider) { super(queueLimit, concurrencyLimit, maxWaitTime, dispatcherThreads, callbackThreads, pollMs, printQueriesFreq, - BufferedRateExecutorType.WRITE, entityService, rateLimitService, serviceInfoProvider, statsFactory, printTenantNames); + BufferedRateExecutorType.WRITE, serviceInfoProvider, rateLimitService, statsFactory, entityService, printTenantNames); } @Scheduled(fixedDelayString = "${cassandra.query.rate_limit_print_interval_ms}") diff --git a/dao/src/main/java/org/thingsboard/server/dao/util/AbstractBufferedRateExecutor.java b/dao/src/main/java/org/thingsboard/server/dao/util/AbstractBufferedRateExecutor.java index 0668c1dd23..cbcf3e81ec 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/util/AbstractBufferedRateExecutor.java +++ b/dao/src/main/java/org/thingsboard/server/dao/util/AbstractBufferedRateExecutor.java @@ -88,12 +88,12 @@ public abstract class AbstractBufferedRateExecutor(queueLimit); @@ -106,7 +106,6 @@ public abstract class AbstractBufferedRateExecutor