Browse Source

Exceeded rate limits notifications; their deduplication; refactoring

pull/8702/head
ViacheslavKlimov 3 years ago
parent
commit
2f560315d1
  1. 6
      application/src/main/java/org/thingsboard/server/config/RateLimitProcessingFilter.java
  2. 2
      application/src/main/java/org/thingsboard/server/controller/AuthController.java
  3. 2
      application/src/main/java/org/thingsboard/server/controller/plugin/TbWebSocketHandler.java
  4. 1
      application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java
  5. 5
      application/src/main/java/org/thingsboard/server/service/install/DefaultSystemDataLoaderService.java
  6. 2
      application/src/main/java/org/thingsboard/server/service/install/SystemDataLoaderService.java
  7. 2
      application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationCenter.java
  8. 2
      application/src/main/java/org/thingsboard/server/service/notification/rule/DefaultNotificationRuleProcessor.java
  9. 71
      application/src/main/java/org/thingsboard/server/service/notification/rule/trigger/RateLimitsTriggerProcessor.java
  10. 2
      application/src/main/java/org/thingsboard/server/service/security/auth/mfa/DefaultTwoFactorAuthService.java
  11. 2
      application/src/main/java/org/thingsboard/server/service/sync/ie/DefaultEntitiesExportImportService.java
  12. 3
      application/src/main/resources/thingsboard.yml
  13. 10
      application/src/test/java/org/thingsboard/server/service/limits/RateLimitServiceTest.java
  14. 2
      application/src/test/java/org/thingsboard/server/service/notification/MockNotificationSettingsService.java
  15. 129
      application/src/test/java/org/thingsboard/server/service/notification/NotificationRuleApiTest.java
  16. 2
      common/dao-api/src/main/java/org/thingsboard/server/dao/notification/NotificationSettingsService.java
  17. 3
      common/dao-api/src/main/java/org/thingsboard/server/dao/notification/NotificationTargetService.java
  18. 70
      common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApi.java
  19. 3
      common/data/src/main/java/org/thingsboard/server/common/data/notification/NotificationType.java
  20. 59
      common/data/src/main/java/org/thingsboard/server/common/data/notification/info/RateLimitsNotificationInfo.java
  21. 1
      common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/trigger/NotificationRuleTriggerConfig.java
  22. 3
      common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/trigger/NotificationRuleTriggerType.java
  23. 45
      common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/trigger/RateLimitsNotificationRuleTriggerConfig.java
  24. 62
      common/message/src/main/java/org/thingsboard/server/common/msg/notification/trigger/RateLimitsTrigger.java
  25. 6
      common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsTransportQueueFactory.java
  26. 6
      common/queue/src/main/java/org/thingsboard/server/queue/provider/InMemoryTbTransportQueueFactory.java
  27. 12
      common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbTransportQueueFactory.java
  28. 7
      common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTransportQueueFactory.java
  29. 7
      common/queue/src/main/java/org/thingsboard/server/queue/provider/RabbitMqTransportQueueFactory.java
  30. 7
      common/queue/src/main/java/org/thingsboard/server/queue/provider/ServiceBusTransportQueueFactory.java
  31. 3
      common/queue/src/main/java/org/thingsboard/server/queue/provider/TbTransportQueueFactory.java
  32. 4
      common/queue/src/main/java/org/thingsboard/server/queue/provider/TbTransportQueueProducerProvider.java
  33. 27
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/DefaultTransportRateLimitService.java
  34. 15
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java
  35. 32
      dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationSettingsService.java
  36. 6
      dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationTargetService.java
  37. 52
      dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotifications.java
  38. 3
      dao/src/main/java/org/thingsboard/server/dao/notification/NotificationTargetDao.java
  39. 6
      dao/src/main/java/org/thingsboard/server/dao/sql/notification/JpaNotificationTargetDao.java
  40. 2
      dao/src/main/java/org/thingsboard/server/dao/util/AbstractBufferedRateExecutor.java
  41. 22
      dao/src/main/java/org/thingsboard/server/dao/util/limits/DefaultRateLimitService.java
  42. 67
      dao/src/main/java/org/thingsboard/server/dao/util/limits/LimitedApi.java
  43. 1
      dao/src/main/java/org/thingsboard/server/dao/util/limits/RateLimitService.java
  44. 7
      msa/vc-executor/src/main/resources/tb-vc-executor.yml
  45. 7
      transport/coap/src/main/resources/tb-coap-transport.yml
  46. 7
      transport/http/src/main/resources/tb-http-transport.yml
  47. 7
      transport/lwm2m/src/main/resources/tb-lwm2m-transport.yml
  48. 7
      transport/mqtt/src/main/resources/tb-mqtt-transport.yml
  49. 7
      transport/snmp/src/main/resources/tb-snmp-transport.yml

6
application/src/main/java/org/thingsboard/server/config/RateLimitProcessingFilter.java

@ -26,7 +26,7 @@ import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.exception.TenantProfileNotFoundException;
import org.thingsboard.server.common.msg.tools.TbRateLimitsException;
import org.thingsboard.server.exception.ThingsboardErrorResponseHandler;
import org.thingsboard.server.dao.util.limits.LimitedApi;
import org.thingsboard.server.common.data.limit.LimitedApi;
import org.thingsboard.server.dao.util.limits.RateLimitService;
import org.thingsboard.server.service.security.model.SecurityUser;
@ -49,7 +49,7 @@ public class RateLimitProcessingFilter extends OncePerRequestFilter {
SecurityUser user = getCurrentUser();
if (user != null && !user.isSystemAdmin()) {
try {
if (!rateLimitService.checkRateLimit(LimitedApi.REST_REQUESTS, user.getTenantId())) {
if (!rateLimitService.checkRateLimit(LimitedApi.REST_REQUESTS_PER_TENANT, user.getTenantId())) {
rateLimitExceeded(EntityType.TENANT, response);
return;
}
@ -60,7 +60,7 @@ public class RateLimitProcessingFilter extends OncePerRequestFilter {
}
if (user.isCustomerUser()) {
if (!rateLimitService.checkRateLimit(LimitedApi.REST_REQUESTS, user.getTenantId(), user.getCustomerId())) {
if (!rateLimitService.checkRateLimit(LimitedApi.REST_REQUESTS_PER_CUSTOMER, user.getTenantId(), user.getCustomerId())) {
rateLimitExceeded(EntityType.CUSTOMER, response);
return;
}

2
application/src/main/java/org/thingsboard/server/controller/AuthController.java

@ -49,7 +49,7 @@ import org.thingsboard.server.common.data.security.model.JwtPair;
import org.thingsboard.server.common.data.security.model.SecuritySettings;
import org.thingsboard.server.common.data.security.model.UserPasswordPolicy;
import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.dao.util.limits.LimitedApi;
import org.thingsboard.server.common.data.limit.LimitedApi;
import org.thingsboard.server.dao.util.limits.RateLimitService;
import org.thingsboard.server.service.security.auth.rest.RestAuthenticationDetails;
import org.thingsboard.server.service.security.model.ActivateUserRequest;

2
application/src/main/java/org/thingsboard/server/controller/plugin/TbWebSocketHandler.java

@ -36,7 +36,7 @@ import org.thingsboard.server.common.data.id.UserId;
import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration;
import org.thingsboard.server.config.WebSocketConfiguration;
import org.thingsboard.server.dao.tenant.TbTenantProfileCache;
import org.thingsboard.server.dao.util.limits.LimitedApi;
import org.thingsboard.server.common.data.limit.LimitedApi;
import org.thingsboard.server.dao.util.limits.RateLimitService;
import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.service.security.model.SecurityUser;

1
application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java

@ -260,6 +260,7 @@ public class ThingsboardInstallService {
case "3.5.1":
log.info("Upgrading ThingsBoard from version 3.5.1 to 3.5.2 ...");
databaseEntitiesUpgradeService.upgradeDatabase("3.5.1");
systemDataLoaderService.updateDefaultNotificationConfigs();
//TODO DON'T FORGET to update switch statement in the CacheCleanupService if you need to clear the cache
break;
default:

5
application/src/main/java/org/thingsboard/server/service/install/DefaultSystemDataLoaderService.java

@ -709,4 +709,9 @@ public class DefaultSystemDataLoaderService implements SystemDataLoaderService {
executor.awaitTermination(Integer.MAX_VALUE, TimeUnit.SECONDS);
}
@Override
public void updateDefaultNotificationConfigs() {
notificationSettingsService.updateDefaultNotificationConfigs(TenantId.SYS_TENANT_ID);
}
}

2
application/src/main/java/org/thingsboard/server/service/install/SystemDataLoaderService.java

@ -41,4 +41,6 @@ public interface SystemDataLoaderService {
void createDefaultNotificationConfigs();
void updateDefaultNotificationConfigs();
}

2
application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationCenter.java

@ -28,6 +28,7 @@ import org.thingsboard.server.common.data.id.NotificationRuleId;
import org.thingsboard.server.common.data.id.NotificationTargetId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.id.UserId;
import org.thingsboard.server.common.data.limit.LimitedApi;
import org.thingsboard.server.common.data.notification.AlreadySentException;
import org.thingsboard.server.common.data.notification.Notification;
import org.thingsboard.server.common.data.notification.NotificationDeliveryMethod;
@ -56,7 +57,6 @@ import org.thingsboard.server.dao.notification.NotificationService;
import org.thingsboard.server.dao.notification.NotificationSettingsService;
import org.thingsboard.server.dao.notification.NotificationTargetService;
import org.thingsboard.server.dao.notification.NotificationTemplateService;
import org.thingsboard.server.dao.util.limits.LimitedApi;
import org.thingsboard.server.dao.util.limits.RateLimitService;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.queue.common.TbProtoQueueMsg;

2
application/src/main/java/org/thingsboard/server/service/notification/rule/DefaultNotificationRuleProcessor.java

@ -32,6 +32,7 @@ import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.NotificationRequestId;
import org.thingsboard.server.common.data.id.NotificationRuleId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.limit.LimitedApi;
import org.thingsboard.server.common.data.notification.NotificationRequest;
import org.thingsboard.server.common.data.notification.NotificationRequestConfig;
import org.thingsboard.server.common.data.notification.NotificationRequestStatus;
@ -46,7 +47,6 @@ import org.thingsboard.server.common.msg.notification.trigger.NotificationRuleTr
import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg;
import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.dao.notification.NotificationRequestService;
import org.thingsboard.server.dao.util.limits.LimitedApi;
import org.thingsboard.server.dao.util.limits.RateLimitService;
import org.thingsboard.server.queue.discovery.PartitionService;
import org.thingsboard.server.service.executors.NotificationExecutorService;

71
application/src/main/java/org/thingsboard/server/service/notification/rule/trigger/RateLimitsTriggerProcessor.java

@ -0,0 +1,71 @@
/**
* Copyright © 2016-2023 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.service.notification.rule.trigger;
import lombok.RequiredArgsConstructor;
import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.notification.info.RateLimitsNotificationInfo;
import org.thingsboard.server.common.data.notification.info.RuleOriginatedNotificationInfo;
import org.thingsboard.server.common.data.notification.rule.trigger.NotificationRuleTriggerType;
import org.thingsboard.server.common.data.notification.rule.trigger.RateLimitsNotificationRuleTriggerConfig;
import org.thingsboard.server.common.data.util.CollectionsUtil;
import org.thingsboard.server.common.msg.notification.trigger.RateLimitsTrigger;
import org.thingsboard.server.dao.entity.EntityService;
import org.thingsboard.server.dao.tenant.TenantService;
import java.util.Optional;
@Service
@RequiredArgsConstructor
public class RateLimitsTriggerProcessor implements NotificationRuleTriggerProcessor<RateLimitsTrigger, RateLimitsNotificationRuleTriggerConfig> {
private final TenantService tenantService;
private final EntityService entityService;
@Override
public boolean matchesFilter(RateLimitsTrigger trigger, RateLimitsNotificationRuleTriggerConfig triggerConfig) {
return trigger.getLimitLevel() != null && trigger.getApi().getLabel() != null &&
CollectionsUtil.emptyOrContains(triggerConfig.getApis(), trigger.getApi());
}
@Override
public RuleOriginatedNotificationInfo constructNotificationInfo(RateLimitsTrigger trigger) {
EntityId limitLevel = trigger.getLimitLevel();
String tenantName = tenantService.findTenantById(trigger.getTenantId()).getName();
String limitLevelEntityName = null;
if (limitLevel instanceof TenantId) {
limitLevelEntityName = tenantName;
} else if (limitLevel != null) {
limitLevelEntityName = Optional.ofNullable(trigger.getLimitLevelEntityName())
.orElseGet(() -> entityService.fetchEntityName(trigger.getTenantId(), limitLevel).orElse(null));
}
return RateLimitsNotificationInfo.builder()
.tenantId(trigger.getTenantId())
.tenantName(tenantName)
.api(trigger.getApi())
.limitLevel(limitLevel)
.limitLevelEntityName(limitLevelEntityName)
.build();
}
@Override
public NotificationRuleTriggerType getTriggerType() {
return NotificationRuleTriggerType.RATE_LIMITS;
}
}

2
application/src/main/java/org/thingsboard/server/service/security/auth/mfa/DefaultTwoFactorAuthService.java

@ -31,7 +31,7 @@ import org.thingsboard.server.common.data.security.model.mfa.account.TwoFaAccoun
import org.thingsboard.server.common.data.security.model.mfa.provider.TwoFaProviderConfig;
import org.thingsboard.server.common.data.security.model.mfa.provider.TwoFaProviderType;
import org.thingsboard.server.dao.user.UserService;
import org.thingsboard.server.dao.util.limits.LimitedApi;
import org.thingsboard.server.common.data.limit.LimitedApi;
import org.thingsboard.server.dao.util.limits.RateLimitService;
import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.service.security.auth.mfa.config.TwoFaConfigManager;

2
application/src/main/java/org/thingsboard/server/service/sync/ie/DefaultEntitiesExportImportService.java

@ -32,7 +32,7 @@ import org.thingsboard.server.common.data.util.ThrowingRunnable;
import org.thingsboard.server.dao.exception.DataValidationException;
import org.thingsboard.server.dao.relation.RelationService;
import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.dao.util.limits.LimitedApi;
import org.thingsboard.server.common.data.limit.LimitedApi;
import org.thingsboard.server.dao.util.limits.RateLimitService;
import org.thingsboard.server.service.entitiy.TbNotificationEntityService;
import org.thingsboard.server.service.sync.ie.exporting.EntityExportService;

3
application/src/main/resources/thingsboard.yml

@ -1266,6 +1266,9 @@ notification_system:
NEW_PLATFORM_VERSION:
# In milliseconds, infinitely by default
deduplication_duration: "${NEW_PLATFORM_VERSION_NOTIFICATION_RULE_DEDUPLICATION_DURATION:0}"
RATE_LIMITS:
# In milliseconds, 4 hours by default
deduplication_duration: "${RATE_LIMITS_NOTIFICATION_RULE_DEDUPLICATION_DURATION:14400000}"
management:
endpoints:

10
application/src/test/java/org/thingsboard/server/service/limits/RateLimitServiceTest.java

@ -24,11 +24,12 @@ import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.NotificationRuleId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.limit.LimitedApi;
import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration;
import org.thingsboard.server.common.data.tenant.profile.TenantProfileData;
import org.thingsboard.server.common.msg.notification.NotificationRuleProcessor;
import org.thingsboard.server.dao.tenant.TbTenantProfileCache;
import org.thingsboard.server.dao.util.limits.DefaultRateLimitService;
import org.thingsboard.server.dao.util.limits.LimitedApi;
import org.thingsboard.server.dao.util.limits.RateLimitService;
import java.util.List;
@ -37,6 +38,7 @@ import java.util.UUID;
import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertTrue;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.reset;
import static org.mockito.Mockito.when;
@ -50,7 +52,7 @@ public class RateLimitServiceTest {
@Before
public void beforeEach() {
tenantProfileCache = Mockito.mock(TbTenantProfileCache.class);
rateLimitService = new DefaultRateLimitService(tenantProfileCache, 60, 100);
rateLimitService = new DefaultRateLimitService(tenantProfileCache, mock(NotificationRuleProcessor.class), 60, 100);
tenantId = new TenantId(UUID.randomUUID());
}
@ -73,14 +75,14 @@ public class RateLimitServiceTest {
LimitedApi.ENTITY_EXPORT,
LimitedApi.ENTITY_IMPORT,
LimitedApi.NOTIFICATION_REQUESTS,
LimitedApi.REST_REQUESTS,
LimitedApi.REST_REQUESTS_PER_CUSTOMER,
LimitedApi.CASSANDRA_QUERIES
)) {
testRateLimits(limitedApi, max, tenantId);
}
CustomerId customerId = new CustomerId(UUID.randomUUID());
testRateLimits(LimitedApi.REST_REQUESTS, max, customerId);
testRateLimits(LimitedApi.REST_REQUESTS_PER_CUSTOMER, max, customerId);
NotificationRuleId notificationRuleId = new NotificationRuleId(UUID.randomUUID());
testRateLimits(LimitedApi.NOTIFICATION_REQUESTS_PER_RULE, max, notificationRuleId);

2
application/src/test/java/org/thingsboard/server/service/notification/MockNotificationSettingsService.java

@ -26,7 +26,7 @@ import org.thingsboard.server.dao.settings.AdminSettingsService;
public class MockNotificationSettingsService extends DefaultNotificationSettingsService {
public MockNotificationSettingsService(AdminSettingsService adminSettingsService) {
super(adminSettingsService, null, null);
super(adminSettingsService, null, null, null);
}
@Override

129
application/src/test/java/org/thingsboard/server/service/notification/NotificationRuleApiTest.java

@ -51,12 +51,15 @@ import org.thingsboard.server.common.data.device.profile.DeviceProfileAlarm;
import org.thingsboard.server.common.data.device.profile.SimpleAlarmConditionSpec;
import org.thingsboard.server.common.data.id.AlarmId;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.limit.LimitedApi;
import org.thingsboard.server.common.data.notification.Notification;
import org.thingsboard.server.common.data.notification.NotificationDeliveryMethod;
import org.thingsboard.server.common.data.notification.NotificationRequest;
import org.thingsboard.server.common.data.notification.NotificationRequestInfo;
import org.thingsboard.server.common.data.notification.NotificationType;
import org.thingsboard.server.common.data.notification.info.AlarmNotificationInfo;
import org.thingsboard.server.common.data.notification.info.RateLimitsNotificationInfo;
import org.thingsboard.server.common.data.notification.rule.DefaultNotificationRuleRecipientsConfig;
import org.thingsboard.server.common.data.notification.rule.EscalatedNotificationRuleRecipientsConfig;
import org.thingsboard.server.common.data.notification.rule.NotificationRule;
@ -70,7 +73,10 @@ import org.thingsboard.server.common.data.notification.rule.trigger.EntitiesLimi
import org.thingsboard.server.common.data.notification.rule.trigger.EntityActionNotificationRuleTriggerConfig;
import org.thingsboard.server.common.data.notification.rule.trigger.NewPlatformVersionNotificationRuleTriggerConfig;
import org.thingsboard.server.common.data.notification.rule.trigger.NotificationRuleTriggerType;
import org.thingsboard.server.common.data.notification.rule.trigger.RateLimitsNotificationRuleTriggerConfig;
import org.thingsboard.server.common.data.notification.targets.NotificationTarget;
import org.thingsboard.server.common.data.notification.targets.platform.AffectedTenantAdministratorsFilter;
import org.thingsboard.server.common.data.notification.targets.platform.SystemAdministratorsFilter;
import org.thingsboard.server.common.data.notification.template.NotificationTemplate;
import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageLink;
@ -82,12 +88,13 @@ import org.thingsboard.server.common.data.rule.RuleChainMetaData;
import org.thingsboard.server.common.data.security.Authority;
import org.thingsboard.server.common.msg.notification.NotificationRuleProcessor;
import org.thingsboard.server.common.msg.notification.trigger.NewPlatformVersionTrigger;
import org.thingsboard.server.common.msg.notification.trigger.RateLimitsTrigger;
import org.thingsboard.server.dao.notification.DefaultNotifications;
import org.thingsboard.server.dao.notification.NotificationRequestService;
import org.thingsboard.server.dao.rule.RuleChainService;
import org.thingsboard.server.dao.service.DaoSqlTest;
import org.thingsboard.server.dao.util.limits.LimitedApi;
import org.thingsboard.server.dao.util.limits.RateLimitService;
import org.thingsboard.server.service.notification.rule.DefaultNotificationRuleProcessor;
import org.thingsboard.server.service.notification.rule.cache.DefaultNotificationRulesCache;
import org.thingsboard.server.service.state.DeviceStateService;
import org.thingsboard.server.service.telemetry.AlarmSubscriptionService;
@ -103,6 +110,7 @@ import java.util.concurrent.Callable;
import java.util.concurrent.TimeUnit;
import java.util.function.BiConsumer;
import java.util.function.Consumer;
import java.util.stream.Collectors;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.offset;
@ -116,7 +124,8 @@ import static org.thingsboard.server.common.data.notification.rule.trigger.Devic
@DaoSqlTest
@TestPropertySource(properties = {
"transport.http.enabled=true"
"transport.http.enabled=true",
"notification_system.rules.trigger_types_configs.RATE_LIMITS.deduplication_duration=10000"
})
public class NotificationRuleApiTest extends AbstractNotificationApiTest {
@ -400,6 +409,60 @@ public class NotificationRuleApiTest extends AbstractNotificationApiTest {
});
}
@Test
public void testNotificationRuleProcessing_exceededRateLimits() throws Exception {
loginSysAdmin();
NotificationTarget sysadmins = createNotificationTarget(new SystemAdministratorsFilter());
NotificationTarget affectedTenantAdmins = createNotificationTarget(new AffectedTenantAdministratorsFilter());
defaultNotifications.create(TenantId.SYS_TENANT_ID, DefaultNotifications.exceededRateLimitsForSysadmin, sysadmins.getId());
defaultNotifications.create(TenantId.SYS_TENANT_ID, DefaultNotifications.exceededRateLimits, affectedTenantAdmins.getId());
defaultNotifications.create(TenantId.SYS_TENANT_ID, DefaultNotifications.exceededPerEntityRateLimits, affectedTenantAdmins.getId());
notificationRulesCache.evict(TenantId.SYS_TENANT_ID);
int n = 10;
updateDefaultTenantProfile(profileConfiguration -> {
profileConfiguration.setTenantEntityExportRateLimit(n + ":600");
profileConfiguration.setCustomerServerRestLimitsConfiguration(n + ":600");
profileConfiguration.setTenantNotificationRequestsPerRuleRateLimit(n + ":600");
profileConfiguration.setTransportDeviceTelemetryMsgRateLimit(n + ":600");
});
loginTenantAdmin();
NotificationRule rule = createNotificationRule(AlarmCommentNotificationRuleTriggerConfig.builder()
.alarmTypes(Set.of("weklfjkwefa"))
.build(), "Test", "Test", createNotificationTarget(tenantAdminUserId).getId());
for (int i = 1; i <= n * 2; i++) {
rateLimitService.checkRateLimit(LimitedApi.ENTITY_EXPORT, tenantId);
rateLimitService.checkRateLimit(LimitedApi.REST_REQUESTS_PER_CUSTOMER, tenantId, customerId);
rateLimitService.checkRateLimit(LimitedApi.NOTIFICATION_REQUESTS_PER_RULE, tenantId, rule.getId());
}
loginTenantAdmin();
List<Notification> notifications = await().atMost(30, TimeUnit.SECONDS)
.until(() -> getMyNotifications(true, 10), list -> list.size() == 3);
assertThat(notifications).allSatisfy(notification -> {
assertThat(notification.getSubject()).isEqualTo("Rate limits exceeded");
});
assertThat(notifications).anySatisfy(notification -> {
assertThat(notification.getText()).isEqualTo("Rate limits for entity version creation exceeded");
});
assertThat(notifications).anySatisfy(notification -> {
assertThat(notification.getText()).isEqualTo("Rate limits for REST API requests per customer " +
"exceeded for 'Customer'");
});
assertThat(notifications).anySatisfy(notification -> {
assertThat(notification.getText()).isEqualTo("Rate limits for notification requests " +
"per rule exceeded for '" + rule.getName() + "'");
});
loginSysAdmin();
notifications = await().atMost(30, TimeUnit.SECONDS)
.until(() -> getMyNotifications(true, 10), list -> list.size() == 1);
assertThat(notifications).allSatisfy(notification -> {
assertThat(notification.getSubject()).isEqualTo("Rate limits exceeded for tenant " + TEST_TENANT_NAME);
});
assertThat(notifications.get(0).getText()).isEqualTo("Rate limits for entity version creation exceeded");
}
@Test
public void testNotificationRuleProcessing_alarmAssignment() throws Exception {
AlarmAssignmentNotificationRuleTriggerConfig triggerConfig = AlarmAssignmentNotificationRuleTriggerConfig.builder()
@ -617,6 +680,68 @@ public class NotificationRuleApiTest extends AbstractNotificationApiTest {
});
}
@Test
public void testNotificationsDeduplication_exceededRateLimits() throws Exception {
RateLimitsNotificationRuleTriggerConfig triggerConfig = new RateLimitsNotificationRuleTriggerConfig();
triggerConfig.setApis(Set.of(LimitedApi.ENTITY_EXPORT, LimitedApi.TRANSPORT_MESSAGES_PER_DEVICE));
loginSysAdmin();
NotificationTarget target = createNotificationTarget(tenantAdminUserId);
NotificationRule rule = createNotificationRule(triggerConfig, "Test 1", "Test", target.getId());
int n = 5;
updateDefaultTenantProfile(profileConfiguration -> {
profileConfiguration.setTenantEntityExportRateLimit(n + ":600");
profileConfiguration.setTransportDeviceTelemetryMsgRateLimit(n + ":800");
});
RateLimitsTrigger expectedTrigger = RateLimitsTrigger.builder()
.tenantId(tenantId)
.api(LimitedApi.ENTITY_EXPORT)
.limitLevel(tenantId)
.build();
assertThat(DefaultNotificationRuleProcessor.getDeduplicationKey(expectedTrigger, rule))
.isEqualTo("RATE_LIMITS:TENANT:" + tenantId + ":ENTITY_EXPORT_" +
target.getId() + ":ENTITY_EXPORT,TRANSPORT_MESSAGES_PER_DEVICE");
loginTenantAdmin();
getWsClient().subscribeForUnreadNotifications(10).waitForReply();
getWsClient().registerWaitForUpdate(2);
Device device = createDevice("Test", "Test");
for (int i = 1; i <= n + 1; i++) {
rateLimitService.checkRateLimit(LimitedApi.ENTITY_EXPORT, tenantId);
doPost("/api/v1/" + device.getName() + "/telemetry", "{\"dp1\":123}", String.class);
}
int expectedNotificationsCount1 = 2;
getWsClient().waitForUpdate(true);
List<Notification> notifications1 = getMyNotifications(true, 10);
assertThat(notifications1).size().isEqualTo(expectedNotificationsCount1);
assertThat(notifications1)
.anyMatch(notification -> ((RateLimitsNotificationInfo) notification.getInfo()).getApi() == LimitedApi.ENTITY_EXPORT)
.anyMatch(notification -> ((RateLimitsNotificationInfo) notification.getInfo()).getApi() == LimitedApi.TRANSPORT_MESSAGES_PER_DEVICE);
getWsClient().registerWaitForUpdate(2);
for (int i = 0; i < 10; i++) {
rateLimitService.checkRateLimit(LimitedApi.ENTITY_EXPORT, tenantId);
doPost("/api/v1/" + device.getName() + "/telemetry", "{\"dp1\":123}", String.class);
}
assertThat(getWsClient().waitForUpdate(5000)).isNull();
int deduplicationDuration = 10000; // configured in TestPropertySource above
await().atLeast(2, TimeUnit.SECONDS)
.atMost(deduplicationDuration, TimeUnit.MILLISECONDS)
.untilAsserted(() -> {
rateLimitService.checkRateLimit(LimitedApi.ENTITY_EXPORT, tenantId);
doPost("/api/v1/" + device.getName() + "/telemetry", "{\"dp1\":123}", String.class);
Map<LimitedApi, Long> notifications2 = getMyNotifications(true, 10).stream()
.map(notification -> (RateLimitsNotificationInfo) notification.getInfo())
.collect(Collectors.groupingBy(RateLimitsNotificationInfo::getApi, Collectors.counting()));
assertThat(notifications2.get(LimitedApi.ENTITY_EXPORT)).isEqualTo(2);
assertThat(notifications2.get(LimitedApi.TRANSPORT_MESSAGES_PER_DEVICE)).isEqualTo(2);
});
}
@Test
public void testNotificationRuleDisabling() throws Exception {
EntityActionNotificationRuleTriggerConfig triggerConfig = new EntityActionNotificationRuleTriggerConfig();

2
common/dao-api/src/main/java/org/thingsboard/server/dao/notification/NotificationSettingsService.java

@ -26,4 +26,6 @@ public interface NotificationSettingsService {
void createDefaultNotificationConfigs(TenantId tenantId);
void updateDefaultNotificationConfigs(TenantId tenantId);
}

3
common/dao-api/src/main/java/org/thingsboard/server/dao/notification/NotificationTargetService.java

@ -23,6 +23,7 @@ import org.thingsboard.server.common.data.notification.NotificationType;
import org.thingsboard.server.common.data.notification.info.RuleOriginatedNotificationInfo;
import org.thingsboard.server.common.data.notification.targets.NotificationTarget;
import org.thingsboard.server.common.data.notification.targets.platform.PlatformUsersNotificationTargetConfig;
import org.thingsboard.server.common.data.notification.targets.platform.UsersFilterType;
import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageLink;
@ -40,6 +41,8 @@ public interface NotificationTargetService {
List<NotificationTarget> findNotificationTargetsByTenantIdAndIds(TenantId tenantId, List<NotificationTargetId> ids);
List<NotificationTarget> findNotificationTargetsByTenantIdAndUsersFilterType(TenantId tenantId, UsersFilterType filterType);
PageData<User> findRecipientsForNotificationTarget(TenantId tenantId, CustomerId customerId, NotificationTargetId targetId, PageLink pageLink);
PageData<User> findRecipientsForNotificationTargetConfig(TenantId tenantId, PlatformUsersNotificationTargetConfig targetConfig, PageLink pageLink);

70
common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApi.java

@ -0,0 +1,70 @@
/**
* Copyright © 2016-2023 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 lombok.Getter;
import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration;
import java.util.Optional;
import java.util.function.Function;
public enum LimitedApi {
ENTITY_EXPORT(DefaultTenantProfileConfiguration::getTenantEntityExportRateLimit, "entity version creation", true),
ENTITY_IMPORT(DefaultTenantProfileConfiguration::getTenantEntityImportRateLimit, "entity version load", true),
NOTIFICATION_REQUESTS(DefaultTenantProfileConfiguration::getTenantNotificationRequestsRateLimit, "notification requests", true),
NOTIFICATION_REQUESTS_PER_RULE(DefaultTenantProfileConfiguration::getTenantNotificationRequestsPerRuleRateLimit, "notification requests per rule", false),
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),
PASSWORD_RESET(false, true),
TWO_FA_VERIFICATION_CODE_SEND(false, true),
TWO_FA_VERIFICATION_CODE_CHECK(false, true),
TRANSPORT_MESSAGES_PER_TENANT("transport messages", true),
TRANSPORT_MESSAGES_PER_DEVICE("transport messages per device", false);
private Function<DefaultTenantProfileConfiguration, String> configExtractor;
@Getter
private final boolean perTenant;
@Getter
private boolean refillRateLimitIntervally;
@Getter
private String label;
LimitedApi(Function<DefaultTenantProfileConfiguration, String> configExtractor, String label, boolean perTenant) {
this.configExtractor = configExtractor;
this.label = label;
this.perTenant = perTenant;
}
LimitedApi(boolean perTenant, boolean refillRateLimitIntervally) {
this.perTenant = perTenant;
this.refillRateLimitIntervally = refillRateLimitIntervally;
}
LimitedApi(String label, boolean perTenant) {
this.label = label;
this.perTenant = perTenant;
}
public String getLimitConfig(DefaultTenantProfileConfiguration profileConfiguration) {
return Optional.ofNullable(configExtractor)
.map(extractor -> extractor.apply(profileConfiguration))
.orElse(null);
}
}

3
common/data/src/main/java/org/thingsboard/server/common/data/notification/NotificationType.java

@ -27,6 +27,7 @@ public enum NotificationType {
NEW_PLATFORM_VERSION,
ENTITIES_LIMIT,
API_USAGE_LIMIT,
RULE_NODE
RULE_NODE,
RATE_LIMITS
}

59
common/data/src/main/java/org/thingsboard/server/common/data/notification/info/RateLimitsNotificationInfo.java

@ -0,0 +1,59 @@
/**
* Copyright © 2016-2023 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.notification.info;
import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
import lombok.NoArgsConstructor;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.limit.LimitedApi;
import java.util.Map;
import static org.thingsboard.server.common.data.util.CollectionsUtil.mapOf;
@Data
@NoArgsConstructor
@AllArgsConstructor
@Builder
public class RateLimitsNotificationInfo implements RuleOriginatedNotificationInfo {
private TenantId tenantId;
private String tenantName;
private LimitedApi api;
private EntityId limitLevel;
private String limitLevelEntityName;
@Override
public Map<String, String> getTemplateData() {
return mapOf(
"api", api.getLabel(),
"limitLevelEntityType", limitLevel != null ? limitLevel.getEntityType().getNormalName() : null,
"limitLevelEntityId", limitLevel != null ? limitLevel.getId().toString() : null,
"limitLevelEntityName", limitLevelEntityName,
"tenantName", tenantName,
"tenantId", tenantId.toString()
);
}
@Override
public TenantId getAffectedTenantId() {
return tenantId;
}
}

1
common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/trigger/NotificationRuleTriggerConfig.java

@ -35,6 +35,7 @@ import java.io.Serializable;
@Type(value = NewPlatformVersionNotificationRuleTriggerConfig.class, name = "NEW_PLATFORM_VERSION"),
@Type(value = EntitiesLimitNotificationRuleTriggerConfig.class, name = "ENTITIES_LIMIT"),
@Type(value = ApiUsageLimitNotificationRuleTriggerConfig.class, name = "API_USAGE_LIMIT"),
@Type(value = RateLimitsNotificationRuleTriggerConfig.class, name = "RATE_LIMITS"),
})
public interface NotificationRuleTriggerConfig extends Serializable {

3
common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/trigger/NotificationRuleTriggerType.java

@ -28,7 +28,8 @@ public enum NotificationRuleTriggerType {
RULE_ENGINE_COMPONENT_LIFECYCLE_EVENT,
NEW_PLATFORM_VERSION(false),
ENTITIES_LIMIT(false),
API_USAGE_LIMIT(false);
API_USAGE_LIMIT(false),
RATE_LIMITS(false);
private final boolean tenantLevel;

45
common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/trigger/RateLimitsNotificationRuleTriggerConfig.java

@ -0,0 +1,45 @@
/**
* Copyright © 2016-2023 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.notification.rule.trigger;
import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
import lombok.NoArgsConstructor;
import org.thingsboard.server.common.data.limit.LimitedApi;
import java.util.Set;
import java.util.stream.Collectors;
@Data
@AllArgsConstructor
@NoArgsConstructor
@Builder
public class RateLimitsNotificationRuleTriggerConfig implements NotificationRuleTriggerConfig {
private Set<LimitedApi> apis;
@Override
public NotificationRuleTriggerType getTriggerType() {
return NotificationRuleTriggerType.RATE_LIMITS;
}
@Override
public String getDeduplicationKey() {
return apis == null ? "#" : apis.stream().sorted().map(Enum::name).collect(Collectors.joining(","));
}
}

62
common/message/src/main/java/org/thingsboard/server/common/msg/notification/trigger/RateLimitsTrigger.java

@ -0,0 +1,62 @@
/**
* Copyright © 2016-2023 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.notification.trigger;
import lombok.Builder;
import lombok.Data;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.limit.LimitedApi;
import org.thingsboard.server.common.data.notification.rule.trigger.NotificationRuleTriggerType;
import java.util.concurrent.TimeUnit;
@Data
@Builder
public class RateLimitsTrigger implements NotificationRuleTrigger {
private final TenantId tenantId;
private final LimitedApi api;
private final EntityId limitLevel;
private final String limitLevelEntityName;
@Override
public NotificationRuleTriggerType getType() {
return NotificationRuleTriggerType.RATE_LIMITS;
}
@Override
public EntityId getOriginatorEntityId() {
return limitLevel != null ? limitLevel : tenantId;
}
@Override
public boolean deduplicate() {
return true;
}
@Override
public String getDeduplicationKey() {
return String.join(":", NotificationRuleTrigger.super.getDeduplicationKey(), api.toString());
}
@Override
public long getDefaultDeduplicationDuration() {
return TimeUnit.HOURS.toMillis(4);
}
}

6
common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsTransportQueueFactory.java

@ -20,6 +20,7 @@ import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.stereotype.Component;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreNotificationMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToTransportMsg;
import org.thingsboard.server.gen.transport.TransportProtos.TransportApiRequestMsg;
@ -110,6 +111,11 @@ public class AwsSqsTransportQueueFactory implements TbTransportQueueFactory {
return new TbAwsSqsProducerTemplate<>(coreAdmin, sqsSettings, coreSettings.getTopic());
}
@Override
public TbQueueProducer<TbProtoQueueMsg<ToCoreNotificationMsg>> createTbCoreNotificationsMsgProducer() {
return new TbAwsSqsProducerTemplate<>(notificationAdmin, sqsSettings, coreSettings.getTopic());
}
@Override
public TbQueueConsumer<TbProtoQueueMsg<ToTransportMsg>> createTransportNotificationsConsumer() {
return new TbAwsSqsConsumerTemplate<>(notificationAdmin, sqsSettings, transportNotificationSettings.getNotificationsTopic() + "_" + serviceInfoProvider.getServiceId(),

6
common/queue/src/main/java/org/thingsboard/server/queue/provider/InMemoryTbTransportQueueFactory.java

@ -20,6 +20,7 @@ import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.stereotype.Component;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreNotificationMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToTransportMsg;
import org.thingsboard.server.gen.transport.TransportProtos.TransportApiRequestMsg;
@ -100,6 +101,11 @@ public class InMemoryTbTransportQueueFactory implements TbTransportQueueFactory
return new InMemoryTbQueueProducer<>(storage, coreSettings.getTopic());
}
@Override
public TbQueueProducer<TbProtoQueueMsg<ToCoreNotificationMsg>> createTbCoreNotificationsMsgProducer() {
return new InMemoryTbQueueProducer<>(storage, coreSettings.getTopic());
}
@Override
public TbQueueConsumer<TbProtoQueueMsg<ToTransportMsg>> createTransportNotificationsConsumer() {
return new InMemoryTbQueueConsumer<>(storage, transportNotificationSettings.getNotificationsTopic() + "." + serviceInfoProvider.getServiceId());

12
common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbTransportQueueFactory.java

@ -18,7 +18,9 @@ package org.thingsboard.server.queue.provider;
import lombok.extern.slf4j.Slf4j;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.stereotype.Component;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreNotificationMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToTransportMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToUsageStatsServiceMsg;
@ -134,6 +136,16 @@ public class KafkaTbTransportQueueFactory implements TbTransportQueueFactory {
return requestBuilder.build();
}
@Override
public TbQueueProducer<TbProtoQueueMsg<ToCoreNotificationMsg>> createTbCoreNotificationsMsgProducer() {
TbKafkaProducerTemplate.TbKafkaProducerTemplateBuilder<TbProtoQueueMsg<ToCoreNotificationMsg>> requestBuilder = TbKafkaProducerTemplate.builder();
requestBuilder.settings(kafkaSettings);
requestBuilder.clientId("transport-node-to-core-notifications-" + serviceInfoProvider.getServiceId());
requestBuilder.defaultTopic(coreSettings.getTopic());
requestBuilder.admin(notificationAdmin);
return requestBuilder.build();
}
@Override
public TbQueueConsumer<TbProtoQueueMsg<ToTransportMsg>> createTransportNotificationsConsumer() {
TbKafkaConsumerTemplate.TbKafkaConsumerTemplateBuilder<TbProtoQueueMsg<ToTransportMsg>> responseBuilder = TbKafkaConsumerTemplate.builder();

7
common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTransportQueueFactory.java

@ -18,7 +18,9 @@ package org.thingsboard.server.queue.provider;
import lombok.extern.slf4j.Slf4j;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.stereotype.Component;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreNotificationMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToTransportMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToUsageStatsServiceMsg;
@ -108,6 +110,11 @@ public class PubSubTransportQueueFactory implements TbTransportQueueFactory {
return new TbPubSubProducerTemplate<>(coreAdmin, pubSubSettings, coreSettings.getTopic());
}
@Override
public TbQueueProducer<TbProtoQueueMsg<ToCoreNotificationMsg>> createTbCoreNotificationsMsgProducer() {
return new TbPubSubProducerTemplate<>(notificationAdmin, pubSubSettings, coreSettings.getTopic());
}
@Override
public TbQueueConsumer<TbProtoQueueMsg<ToTransportMsg>> createTransportNotificationsConsumer() {
return new TbPubSubConsumerTemplate<>(notificationAdmin, pubSubSettings,

7
common/queue/src/main/java/org/thingsboard/server/queue/provider/RabbitMqTransportQueueFactory.java

@ -18,7 +18,9 @@ package org.thingsboard.server.queue.provider;
import lombok.extern.slf4j.Slf4j;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.stereotype.Component;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreNotificationMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToTransportMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToUsageStatsServiceMsg;
@ -110,6 +112,11 @@ public class RabbitMqTransportQueueFactory implements TbTransportQueueFactory {
return new TbRabbitMqProducerTemplate<>(coreAdmin, rabbitMqSettings, coreSettings.getTopic());
}
@Override
public TbQueueProducer<TbProtoQueueMsg<ToCoreNotificationMsg>> createTbCoreNotificationsMsgProducer() {
return new TbRabbitMqProducerTemplate<>(notificationAdmin, rabbitMqSettings, coreSettings.getTopic());
}
@Override
public TbQueueConsumer<TbProtoQueueMsg<ToTransportMsg>> createTransportNotificationsConsumer() {
return new TbRabbitMqConsumerTemplate<>(notificationAdmin, rabbitMqSettings, transportNotificationSettings.getNotificationsTopic() + "." + serviceInfoProvider.getServiceId(),

7
common/queue/src/main/java/org/thingsboard/server/queue/provider/ServiceBusTransportQueueFactory.java

@ -18,7 +18,9 @@ package org.thingsboard.server.queue.provider;
import lombok.extern.slf4j.Slf4j;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.stereotype.Component;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreNotificationMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToTransportMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToUsageStatsServiceMsg;
@ -110,6 +112,11 @@ public class ServiceBusTransportQueueFactory implements TbTransportQueueFactory
return new TbServiceBusProducerTemplate<>(coreAdmin, serviceBusSettings, coreSettings.getTopic());
}
@Override
public TbQueueProducer<TbProtoQueueMsg<ToCoreNotificationMsg>> createTbCoreNotificationsMsgProducer() {
return new TbServiceBusProducerTemplate<>(notificationAdmin, serviceBusSettings, coreSettings.getTopic());
}
@Override
public TbQueueConsumer<TbProtoQueueMsg<ToTransportMsg>> createTransportNotificationsConsumer() {
return new TbServiceBusConsumerTemplate<>(notificationAdmin, serviceBusSettings,

3
common/queue/src/main/java/org/thingsboard/server/queue/provider/TbTransportQueueFactory.java

@ -16,6 +16,7 @@
package org.thingsboard.server.queue.provider;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreNotificationMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToTransportMsg;
import org.thingsboard.server.gen.transport.TransportProtos.TransportApiRequestMsg;
@ -33,6 +34,8 @@ public interface TbTransportQueueFactory extends TbUsageStatsClientQueueFactory
TbQueueProducer<TbProtoQueueMsg<ToCoreMsg>> createTbCoreMsgProducer();
TbQueueProducer<TbProtoQueueMsg<ToCoreNotificationMsg>> createTbCoreNotificationsMsgProducer();
TbQueueConsumer<TbProtoQueueMsg<ToTransportMsg>> createTransportNotificationsConsumer();
}

4
common/queue/src/main/java/org/thingsboard/server/queue/provider/TbTransportQueueProducerProvider.java

@ -36,6 +36,7 @@ public class TbTransportQueueProducerProvider implements TbQueueProducerProvider
private final TbTransportQueueFactory tbQueueProvider;
private TbQueueProducer<TbProtoQueueMsg<ToRuleEngineMsg>> toRuleEngine;
private TbQueueProducer<TbProtoQueueMsg<ToCoreMsg>> toTbCore;
private TbQueueProducer<TbProtoQueueMsg<ToCoreNotificationMsg>> toTbCoreNotifications;
private TbQueueProducer<TbProtoQueueMsg<ToUsageStatsServiceMsg>> toUsageStats;
public TbTransportQueueProducerProvider(TbTransportQueueFactory tbQueueProvider) {
@ -47,6 +48,7 @@ public class TbTransportQueueProducerProvider implements TbQueueProducerProvider
this.toTbCore = tbQueueProvider.createTbCoreMsgProducer();
this.toRuleEngine = tbQueueProvider.createRuleEngineMsgProducer();
this.toUsageStats = tbQueueProvider.createToUsageStatsServiceMsgProducer();
this.toTbCoreNotifications = tbQueueProvider.createTbCoreNotificationsMsgProducer();
}
@Override
@ -71,7 +73,7 @@ public class TbTransportQueueProducerProvider implements TbQueueProducerProvider
@Override
public TbQueueProducer<TbProtoQueueMsg<ToCoreNotificationMsg>> getTbCoreNotificationsMsgProducer() {
throw new RuntimeException("Not Implemented! Should not be used by Transport!");
return toTbCoreNotifications;
}
@Override

27
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/DefaultTransportRateLimitService.java

@ -78,11 +78,11 @@ public class DefaultTransportRateLimitService implements TransportRateLimitServi
return null;
}
private boolean checkEntityRateLimit(int dataPoints, EntityTransportRateLimits tenantLimits) {
private boolean checkEntityRateLimit(int dataPoints, EntityTransportRateLimits limits) {
if (dataPoints > 0) {
return tenantLimits.getTelemetryMsgRateLimit().tryConsume() && tenantLimits.getTelemetryDataPointsRateLimit().tryConsume(dataPoints);
return limits.getTelemetryMsgRateLimit().tryConsume() && limits.getTelemetryDataPointsRateLimit().tryConsume(dataPoints);
} else {
return tenantLimits.getRegularMsgRateLimit().tryConsume();
return limits.getRegularMsgRateLimit().tryConsume();
}
}
@ -241,7 +241,7 @@ public class DefaultTransportRateLimitService implements TransportRateLimitServi
} else {
TransportRateLimit regularMsgRateLimit = newLimit(tenant ? profile.getTransportTenantMsgRateLimit() : profile.getTransportDeviceMsgRateLimit());
TransportRateLimit telemetryMsgRateLimit = newLimit(tenant ? profile.getTransportTenantTelemetryMsgRateLimit() : profile.getTransportDeviceTelemetryMsgRateLimit());
TransportRateLimit telemetryDpRateLimit = newLimit(tenant ? profile.getTransportTenantTelemetryDataPointsRateLimit() : profile.getTransportTenantTelemetryDataPointsRateLimit());
TransportRateLimit telemetryDpRateLimit = newLimit(tenant ? profile.getTransportTenantTelemetryDataPointsRateLimit() : profile.getTransportDeviceTelemetryDataPointsRateLimit());
return new EntityTransportRateLimits(regularMsgRateLimit, telemetryMsgRateLimit, telemetryDpRateLimit);
}
}
@ -251,21 +251,16 @@ public class DefaultTransportRateLimitService implements TransportRateLimitServi
}
private EntityTransportRateLimits getTenantRateLimits(TenantId tenantId) {
EntityTransportRateLimits limits = perTenantLimits.get(tenantId);
if (limits == null) {
limits = createRateLimits(tenantProfileCache.get(tenantId), true);
perTenantLimits.put(tenantId, limits);
}
return limits;
return perTenantLimits.computeIfAbsent(tenantId, k -> {
return createRateLimits(tenantProfileCache.get(tenantId), true);
});
}
private EntityTransportRateLimits getDeviceRateLimits(TenantId tenantId, DeviceId deviceId) {
EntityTransportRateLimits limits = perDeviceLimits.get(deviceId);
if (limits == null) {
limits = createRateLimits(tenantProfileCache.get(tenantId), false);
perDeviceLimits.put(deviceId, limits);
return perDeviceLimits.computeIfAbsent(deviceId, k -> {
EntityTransportRateLimits limits = createRateLimits(tenantProfileCache.get(tenantId), false);
tenantDevices.computeIfAbsent(tenantId, id -> ConcurrentHashMap.newKeySet()).add(deviceId);
}
return limits;
return limits;
});
}
}

15
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java

@ -48,9 +48,12 @@ import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.RuleChainId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.id.TenantProfileId;
import org.thingsboard.server.common.data.limit.LimitedApi;
import org.thingsboard.server.common.data.rpc.RpcStatus;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgMetaData;
import org.thingsboard.server.common.msg.notification.NotificationRuleProcessor;
import org.thingsboard.server.common.msg.notification.trigger.RateLimitsTrigger;
import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
import org.thingsboard.server.common.msg.session.SessionMsgType;
@ -177,6 +180,7 @@ public class DefaultTransportService implements TransportService {
private final SchedulerComponent scheduler;
private final ApplicationEventPublisher eventPublisher;
private final TransportResourceCache transportResourceCache;
private final NotificationRuleProcessor notificationRuleProcessor;
protected TbQueueRequestTemplate<TbProtoQueueMsg<TransportApiRequestMsg>, TbProtoQueueMsg<TransportApiResponseMsg>> transportApiRequestTemplate;
protected TbQueueProducer<TbProtoQueueMsg<ToRuleEngineMsg>> ruleEngineMsgProducer;
@ -206,7 +210,7 @@ public class DefaultTransportService implements TransportService {
TransportTenantProfileCache tenantProfileCache,
TransportRateLimitService rateLimitService,
DataDecodingEncodingService dataDecodingEncodingService, SchedulerComponent scheduler, TransportResourceCache transportResourceCache,
ApplicationEventPublisher eventPublisher) {
ApplicationEventPublisher eventPublisher, NotificationRuleProcessor notificationRuleProcessor) {
this.partitionService = partitionService;
this.serviceInfoProvider = serviceInfoProvider;
this.queueProvider = queueProvider;
@ -220,6 +224,7 @@ public class DefaultTransportService implements TransportService {
this.scheduler = scheduler;
this.transportResourceCache = transportResourceCache;
this.eventPublisher = eventPublisher;
this.notificationRuleProcessor = notificationRuleProcessor;
}
@PostConstruct
@ -877,6 +882,14 @@ public class DefaultTransportService implements TransportService {
if (callback != null) {
callback.onError(new TbRateLimitsException(rateLimitedEntityType));
}
if (rateLimitedEntityType == EntityType.DEVICE || rateLimitedEntityType == EntityType.TENANT) {
notificationRuleProcessor.process(RateLimitsTrigger.builder()
.tenantId(tenantId)
.api(rateLimitedEntityType == EntityType.DEVICE ? LimitedApi.TRANSPORT_MESSAGES_PER_DEVICE : LimitedApi.TRANSPORT_MESSAGES_PER_TENANT)
.limitLevel(rateLimitedEntityType == EntityType.DEVICE ? deviceId : tenantId)
.limitLevelEntityName(rateLimitedEntityType == EntityType.DEVICE ? sessionInfo.getDeviceName() : null)
.build());
}
return false;
}
}

32
dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationSettingsService.java

@ -25,6 +25,7 @@ import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.common.data.AdminSettings;
import org.thingsboard.server.common.data.CacheConstants;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.notification.NotificationType;
import org.thingsboard.server.common.data.notification.settings.NotificationSettings;
import org.thingsboard.server.common.data.notification.targets.NotificationTarget;
import org.thingsboard.server.common.data.notification.targets.platform.AffectedTenantAdministratorsFilter;
@ -35,9 +36,12 @@ import org.thingsboard.server.common.data.notification.targets.platform.Platform
import org.thingsboard.server.common.data.notification.targets.platform.SystemAdministratorsFilter;
import org.thingsboard.server.common.data.notification.targets.platform.TenantAdministratorsFilter;
import org.thingsboard.server.common.data.notification.targets.platform.UsersFilter;
import org.thingsboard.server.common.data.notification.targets.platform.UsersFilterType;
import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.dao.settings.AdminSettingsService;
import java.util.Collections;
import java.util.List;
import java.util.Optional;
@Service
@ -46,6 +50,7 @@ public class DefaultNotificationSettingsService implements NotificationSettingsS
private final AdminSettingsService adminSettingsService;
private final NotificationTargetService notificationTargetService;
private final NotificationTemplateService notificationTemplateService;
private final DefaultNotifications defaultNotifications;
private static final String SETTINGS_KEY = "notifications";
@ -98,6 +103,10 @@ public class DefaultNotificationSettingsService implements NotificationSettingsS
defaultNotifications.create(tenantId, DefaultNotifications.apiFeatureDisabledForSysadmin, sysAdmins.getId());
defaultNotifications.create(tenantId, DefaultNotifications.apiFeatureDisabledForTenant, affectedTenantAdmins.getId());
defaultNotifications.create(tenantId, DefaultNotifications.exceededRateLimits, affectedTenantAdmins.getId());
defaultNotifications.create(tenantId, DefaultNotifications.exceededPerEntityRateLimits, affectedTenantAdmins.getId());
defaultNotifications.create(tenantId, DefaultNotifications.exceededRateLimitsForSysadmin, sysAdmins.getId());
defaultNotifications.create(tenantId, DefaultNotifications.newPlatformVersion, sysAdmins.getId());
return;
}
@ -116,6 +125,29 @@ public class DefaultNotificationSettingsService implements NotificationSettingsS
defaultNotifications.create(tenantId, DefaultNotifications.ruleEngineComponentLifecycleFailure, tenantAdmins.getId());
}
@Override
public void updateDefaultNotificationConfigs(TenantId tenantId) {
if (tenantId.isSysTenantId()) {
if (notificationTemplateService.findNotificationTemplatesByTenantIdAndNotificationTypes(tenantId,
List.of(NotificationType.RATE_LIMITS), new PageLink(1)).getTotalElements() > 0) {
return;
}
NotificationTarget sysAdmins = notificationTargetService.findNotificationTargetsByTenantIdAndUsersFilterType(tenantId, UsersFilterType.SYSTEM_ADMINISTRATORS).stream()
.findFirst().orElseGet(() -> {
return createTarget(tenantId, "System administrators", new SystemAdministratorsFilter(), "All system administrators");
});
NotificationTarget affectedTenantAdmins = notificationTargetService.findNotificationTargetsByTenantIdAndUsersFilterType(tenantId, UsersFilterType.AFFECTED_TENANT_ADMINISTRATORS).stream()
.findFirst().orElseGet(() -> {
return createTarget(tenantId, "Affected tenant's administrators", new AffectedTenantAdministratorsFilter(), "");
});
defaultNotifications.create(tenantId, DefaultNotifications.exceededRateLimits, affectedTenantAdmins.getId());
defaultNotifications.create(tenantId, DefaultNotifications.exceededPerEntityRateLimits, affectedTenantAdmins.getId());
defaultNotifications.create(tenantId, DefaultNotifications.exceededRateLimitsForSysadmin, sysAdmins.getId());
}
}
private NotificationTarget createTarget(TenantId tenantId, String name, UsersFilter filter, String description) {
NotificationTarget target = new NotificationTarget();
target.setTenantId(tenantId);

6
dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationTargetService.java

@ -37,6 +37,7 @@ import org.thingsboard.server.common.data.notification.targets.platform.Platform
import org.thingsboard.server.common.data.notification.targets.platform.TenantAdministratorsFilter;
import org.thingsboard.server.common.data.notification.targets.platform.UserListFilter;
import org.thingsboard.server.common.data.notification.targets.platform.UsersFilter;
import org.thingsboard.server.common.data.notification.targets.platform.UsersFilterType;
import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.dao.entity.AbstractEntityService;
@ -93,6 +94,11 @@ public class DefaultNotificationTargetService extends AbstractEntityService impl
return notificationTargetDao.findByTenantIdAndIds(tenantId, ids);
}
@Override
public List<NotificationTarget> findNotificationTargetsByTenantIdAndUsersFilterType(TenantId tenantId, UsersFilterType filterType) {
return notificationTargetDao.findByTenantIdAndUsersFilterType(tenantId, filterType);
}
@Override
public PageData<User> findRecipientsForNotificationTarget(TenantId tenantId, CustomerId customerId, NotificationTargetId targetId, PageLink pageLink) {
NotificationTarget notificationTarget = findNotificationTargetById(tenantId, targetId);

52
dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotifications.java

@ -26,6 +26,7 @@ import org.thingsboard.server.common.data.alarm.AlarmSearchStatus;
import org.thingsboard.server.common.data.id.NotificationTargetId;
import org.thingsboard.server.common.data.id.NotificationTemplateId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.limit.LimitedApi;
import org.thingsboard.server.common.data.notification.NotificationDeliveryMethod;
import org.thingsboard.server.common.data.notification.NotificationType;
import org.thingsboard.server.common.data.notification.rule.DefaultNotificationRuleRecipientsConfig;
@ -44,16 +45,20 @@ import org.thingsboard.server.common.data.notification.rule.trigger.EntityAction
import org.thingsboard.server.common.data.notification.rule.trigger.NewPlatformVersionNotificationRuleTriggerConfig;
import org.thingsboard.server.common.data.notification.rule.trigger.NotificationRuleTriggerConfig;
import org.thingsboard.server.common.data.notification.rule.trigger.NotificationRuleTriggerType;
import org.thingsboard.server.common.data.notification.rule.trigger.RateLimitsNotificationRuleTriggerConfig;
import org.thingsboard.server.common.data.notification.rule.trigger.RuleEngineComponentLifecycleEventNotificationRuleTriggerConfig;
import org.thingsboard.server.common.data.notification.template.NotificationTemplate;
import org.thingsboard.server.common.data.notification.template.NotificationTemplateConfig;
import org.thingsboard.server.common.data.notification.template.WebDeliveryMethodNotificationTemplate;
import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent;
import java.util.Arrays;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.stream.Collectors;
import static java.util.function.Predicate.not;
import static org.thingsboard.common.util.JacksonUtil.newObjectNode;
import static org.thingsboard.server.dao.DaoUtil.toUUIDs;
@ -66,6 +71,7 @@ public class DefaultNotifications {
.subject("Infrastructure maintenance")
.text("Maintenance work is scheduled for tomorrow (7:00 a.m. - 9:00 a.m. UTC)")
.build();
public static final DefaultNotification entitiesLimitForSysadmin = DefaultNotification.builder()
.name("Entities count limit notification for sysadmin")
.type(NotificationType.ENTITIES_LIMIT)
@ -88,6 +94,7 @@ public class DefaultNotifications {
.description("Send notification to tenant admins when count of entities of some type reached 80% threshold of the limit")
.build())
.build();
public static final DefaultNotification apiFeatureWarningForSysadmin = DefaultNotification.builder()
.name("API feature warning notification for sysadmin")
.type(NotificationType.API_USAGE_LIMIT)
@ -134,6 +141,51 @@ public class DefaultNotifications {
.description("Send notification to tenant admins when API feature is disabled")
.build())
.build();
public static final DefaultNotification exceededRateLimits = DefaultNotification.builder()
.name("Exceeded per-tenant rate limits notification for tenant")
.type(NotificationType.RATE_LIMITS)
.subject("Rate limits exceeded")
.text("Rate limits for ${api} exceeded")
.icon("block").color("#e91a1a")
.rule(DefaultRule.builder()
.name("Per-tenant rate limits exceeded")
.triggerConfig(RateLimitsNotificationRuleTriggerConfig.builder()
.apis(Arrays.stream(LimitedApi.values())
.filter(LimitedApi::isPerTenant)
.filter(api -> api.getLabel() != null)
.collect(Collectors.toSet()))
.build())
.description("Send notification to tenant admins when some per-tenant rate limit is exceeded")
.build())
.build();
public static final DefaultNotification exceededPerEntityRateLimits = DefaultNotification.builder()
.name("Exceeded per-entity rate limits notification for tenant")
.type(NotificationType.RATE_LIMITS)
.subject("Rate limits exceeded")
.text("Rate limits for ${api} exceeded for '${limitLevelEntityName}'")
.icon("block").color("#e91a1a")
.rule(DefaultRule.builder()
.name("Per-entity rate limits exceeded")
.triggerConfig(RateLimitsNotificationRuleTriggerConfig.builder()
.apis(Arrays.stream(LimitedApi.values())
.filter(not(LimitedApi::isPerTenant))
.filter(api -> api.getLabel() != null)
.collect(Collectors.toSet()))
.build())
.description("Send notification to tenant admins when some per-entity rate limit is exceeded for an entity")
.build())
.build();
public static final DefaultNotification exceededRateLimitsForSysadmin = exceededRateLimits.toBuilder()
.name("Exceeded per-tenant rate limits notification for sysadmin")
.subject("Rate limits exceeded for tenant ${tenantName}")
.button("Go to tenant").link("/tenants/${tenantId}")
.rule(exceededRateLimits.getRule().toBuilder()
.name("Per-tenant rate limits exceeded (sysadmin)")
.description("Send notification to system admins when a tenant exceeds some per-tenant rate limit")
.build())
.build();
public static final DefaultNotification newPlatformVersion = DefaultNotification.builder()
.name("New platform version notification")
.type(NotificationType.NEW_PLATFORM_VERSION)

3
dao/src/main/java/org/thingsboard/server/dao/notification/NotificationTargetDao.java

@ -19,6 +19,7 @@ import org.thingsboard.server.common.data.id.NotificationTargetId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.notification.NotificationType;
import org.thingsboard.server.common.data.notification.targets.NotificationTarget;
import org.thingsboard.server.common.data.notification.targets.platform.UsersFilterType;
import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.dao.Dao;
@ -35,6 +36,8 @@ public interface NotificationTargetDao extends Dao<NotificationTarget>, TenantEn
List<NotificationTarget> findByTenantIdAndIds(TenantId tenantId, List<NotificationTargetId> ids);
List<NotificationTarget> findByTenantIdAndUsersFilterType(TenantId tenantId, UsersFilterType filterType);
void removeByTenantId(TenantId tenantId);
}

6
dao/src/main/java/org/thingsboard/server/dao/sql/notification/JpaNotificationTargetDao.java

@ -66,6 +66,12 @@ public class JpaNotificationTargetDao extends JpaAbstractDao<NotificationTargetE
return DaoUtil.convertDataList(notificationTargetRepository.findByTenantIdAndIdIn(tenantId.getId(), DaoUtil.toUUIDs(ids)));
}
@Override
public List<NotificationTarget> findByTenantIdAndUsersFilterType(TenantId tenantId, UsersFilterType filterType) {
return DaoUtil.convertDataList(notificationTargetRepository.findByTenantIdAndSearchTextAndUsersFilterTypeIfPresent(tenantId.getId(), "",
List.of(filterType.name()), DaoUtil.toPageable(new PageLink(Integer.MAX_VALUE))).getContent());
}
@Override
public void removeByTenantId(TenantId tenantId) {
notificationTargetRepository.deleteByTenantId(tenantId.getId());

2
dao/src/main/java/org/thingsboard/server/dao/util/AbstractBufferedRateExecutor.java

@ -37,7 +37,7 @@ 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.dao.util.limits.LimitedApi;
import org.thingsboard.server.common.data.limit.LimitedApi;
import org.thingsboard.server.dao.util.limits.RateLimitService;
import javax.annotation.Nullable;

22
dao/src/main/java/org/thingsboard/server/dao/util/limits/DefaultRateLimitService.java

@ -20,11 +20,16 @@ import com.github.benmanes.caffeine.cache.Caffeine;
import lombok.Data;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Lazy;
import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.exception.TenantProfileNotFoundException;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.limit.LimitedApi;
import org.thingsboard.server.common.msg.notification.NotificationRuleProcessor;
import org.thingsboard.server.common.msg.notification.trigger.RateLimitsTrigger;
import org.thingsboard.server.common.msg.tools.TbRateLimits;
import org.thingsboard.server.dao.tenant.TbTenantProfileCache;
@ -35,11 +40,14 @@ import java.util.concurrent.TimeUnit;
public class DefaultRateLimitService implements RateLimitService {
private final TbTenantProfileCache tenantProfileCache;
private final NotificationRuleProcessor notificationRuleProcessor;
public DefaultRateLimitService(TbTenantProfileCache tenantProfileCache,
@Lazy NotificationRuleProcessor notificationRuleProcessor,
@Value("${cache.rateLimits.timeToLiveInMinutes:120}") int rateLimitsTtl,
@Value("${cache.rateLimits.maxSize:200000}") int rateLimitsCacheMaxSize) {
this.tenantProfileCache = tenantProfileCache;
this.notificationRuleProcessor = notificationRuleProcessor;
this.rateLimits = Caffeine.newBuilder()
.expireAfterAccess(rateLimitsTtl, TimeUnit.MINUTES)
.maximumSize(rateLimitsCacheMaxSize)
@ -64,9 +72,17 @@ public class DefaultRateLimitService implements RateLimitService {
}
String rateLimitConfig = tenantProfile.getProfileConfiguration()
.map(profileConfiguration -> api.getLimitConfig(profileConfiguration, level))
.orElse(null);
return checkRateLimit(api, level, rateLimitConfig);
.map(api::getLimitConfig).orElse(null);
boolean success = checkRateLimit(api, level, rateLimitConfig);
if (!success) {
notificationRuleProcessor.process(RateLimitsTrigger.builder()
.tenantId(tenantId)
.api(api)
.limitLevel(level instanceof EntityId ? (EntityId) level : tenantId)
.limitLevelEntityName(null)
.build());
}
return success;
}
@Override

67
dao/src/main/java/org/thingsboard/server/dao/util/limits/LimitedApi.java

@ -1,67 +0,0 @@
/**
* Copyright © 2016-2023 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.limits;
import lombok.Getter;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration;
import java.util.function.BiFunction;
import java.util.function.Function;
public enum LimitedApi {
ENTITY_EXPORT(DefaultTenantProfileConfiguration::getTenantEntityExportRateLimit),
ENTITY_IMPORT(DefaultTenantProfileConfiguration::getTenantEntityImportRateLimit),
NOTIFICATION_REQUESTS(DefaultTenantProfileConfiguration::getTenantNotificationRequestsRateLimit),
NOTIFICATION_REQUESTS_PER_RULE(DefaultTenantProfileConfiguration::getTenantNotificationRequestsPerRuleRateLimit),
REST_REQUESTS((profileConfiguration, level) -> ((EntityId) level).getEntityType() == EntityType.TENANT ?
profileConfiguration.getTenantServerRestLimitsConfiguration() :
profileConfiguration.getCustomerServerRestLimitsConfiguration()),
WS_UPDATES_PER_SESSION(DefaultTenantProfileConfiguration::getWsUpdatesPerSessionRateLimit),
CASSANDRA_QUERIES(DefaultTenantProfileConfiguration::getCassandraQueryTenantRateLimitsConfiguration),
PASSWORD_RESET(true),
TWO_FA_VERIFICATION_CODE_SEND(true),
TWO_FA_VERIFICATION_CODE_CHECK(true);
private final BiFunction<DefaultTenantProfileConfiguration, Object, String> configExtractor;
@Getter
private final boolean refillRateLimitIntervally;
LimitedApi(Function<DefaultTenantProfileConfiguration, String> configExtractor) {
this((profileConfiguration, level) -> configExtractor.apply(profileConfiguration));
}
LimitedApi(BiFunction<DefaultTenantProfileConfiguration, Object, String> configExtractor) {
this.configExtractor = configExtractor;
this.refillRateLimitIntervally = false;
}
LimitedApi(boolean refillRateLimitIntervally) {
this.configExtractor = null;
this.refillRateLimitIntervally = refillRateLimitIntervally;
}
public String getLimitConfig(DefaultTenantProfileConfiguration profileConfiguration, Object level) {
if (configExtractor != null) {
return configExtractor.apply(profileConfiguration, level);
} else {
throw new IllegalArgumentException("No tenant profile config for " + name() + " rate limits");
}
}
}

1
dao/src/main/java/org/thingsboard/server/dao/util/limits/RateLimitService.java

@ -16,6 +16,7 @@
package org.thingsboard.server.dao.util.limits;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.limit.LimitedApi;
public interface RateLimitService {

7
msa/vc-executor/src/main/resources/tb-vc-executor.yml

@ -203,3 +203,10 @@ service:
type: "${TB_SERVICE_TYPE:tb-vc-executor}"
# Unique id for this service (autogenerated if empty)
id: "${TB_SERVICE_ID:}"
notification_system:
rules:
trigger_types_configs:
RATE_LIMITS:
# In milliseconds, 4 hours by default
deduplication_duration: "${RATE_LIMITS_NOTIFICATION_RULE_DEDUPLICATION_DURATION:14400000}"

7
transport/coap/src/main/resources/tb-coap-transport.yml

@ -302,3 +302,10 @@ management:
exposure:
# Expose metrics endpoint (use value 'prometheus' to enable prometheus metrics).
include: '${METRICS_ENDPOINTS_EXPOSE:info}'
notification_system:
rules:
trigger_types_configs:
RATE_LIMITS:
# In milliseconds, 4 hours by default
deduplication_duration: "${RATE_LIMITS_NOTIFICATION_RULE_DEDUPLICATION_DURATION:14400000}"

7
transport/http/src/main/resources/tb-http-transport.yml

@ -287,3 +287,10 @@ management:
exposure:
# Expose metrics endpoint (use value 'prometheus' to enable prometheus metrics).
include: '${METRICS_ENDPOINTS_EXPOSE:info}'
notification_system:
rules:
trigger_types_configs:
RATE_LIMITS:
# In milliseconds, 4 hours by default
deduplication_duration: "${RATE_LIMITS_NOTIFICATION_RULE_DEDUPLICATION_DURATION:14400000}"

7
transport/lwm2m/src/main/resources/tb-lwm2m-transport.yml

@ -369,3 +369,10 @@ management:
exposure:
# Expose metrics endpoint (use value 'prometheus' to enable prometheus metrics).
include: '${METRICS_ENDPOINTS_EXPOSE:info}'
notification_system:
rules:
trigger_types_configs:
RATE_LIMITS:
# In milliseconds, 4 hours by default
deduplication_duration: "${RATE_LIMITS_NOTIFICATION_RULE_DEDUPLICATION_DURATION:14400000}"

7
transport/mqtt/src/main/resources/tb-mqtt-transport.yml

@ -317,3 +317,10 @@ management:
exposure:
# Expose metrics endpoint (use value 'prometheus' to enable prometheus metrics).
include: '${METRICS_ENDPOINTS_EXPOSE:info}'
notification_system:
rules:
trigger_types_configs:
RATE_LIMITS:
# In milliseconds, 4 hours by default
deduplication_duration: "${RATE_LIMITS_NOTIFICATION_RULE_DEDUPLICATION_DURATION:14400000}"

7
transport/snmp/src/main/resources/tb-snmp-transport.yml

@ -267,3 +267,10 @@ management:
exposure:
# Expose metrics endpoint (use value 'prometheus' to enable prometheus metrics).
include: '${METRICS_ENDPOINTS_EXPOSE:info}'
notification_system:
rules:
trigger_types_configs:
RATE_LIMITS:
# In milliseconds, 4 hours by default
deduplication_duration: "${RATE_LIMITS_NOTIFICATION_RULE_DEDUPLICATION_DURATION:14400000}"

Loading…
Cancel
Save