From 2f560315d16b759fd985a979a4eba84c829a4c00 Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Fri, 2 Jun 2023 15:09:09 +0300 Subject: [PATCH] Exceeded rate limits notifications; their deduplication; refactoring --- .../config/RateLimitProcessingFilter.java | 6 +- .../server/controller/AuthController.java | 2 +- .../controller/plugin/TbWebSocketHandler.java | 2 +- .../install/ThingsboardInstallService.java | 1 + .../DefaultSystemDataLoaderService.java | 5 + .../install/SystemDataLoaderService.java | 2 + .../DefaultNotificationCenter.java | 2 +- .../DefaultNotificationRuleProcessor.java | 2 +- .../trigger/RateLimitsTriggerProcessor.java | 71 ++++++++++ .../auth/mfa/DefaultTwoFactorAuthService.java | 2 +- .../DefaultEntitiesExportImportService.java | 2 +- .../src/main/resources/thingsboard.yml | 3 + .../service/limits/RateLimitServiceTest.java | 10 +- .../MockNotificationSettingsService.java | 2 +- .../notification/NotificationRuleApiTest.java | 129 +++++++++++++++++- .../NotificationSettingsService.java | 2 + .../NotificationTargetService.java | 3 + .../server/common/data/limit/LimitedApi.java | 70 ++++++++++ .../data/notification/NotificationType.java | 3 +- .../info/RateLimitsNotificationInfo.java | 59 ++++++++ .../NotificationRuleTriggerConfig.java | 1 + .../trigger/NotificationRuleTriggerType.java | 3 +- ...teLimitsNotificationRuleTriggerConfig.java | 45 ++++++ .../trigger/RateLimitsTrigger.java | 62 +++++++++ .../provider/AwsSqsTransportQueueFactory.java | 6 + .../InMemoryTbTransportQueueFactory.java | 6 + .../KafkaTbTransportQueueFactory.java | 12 ++ .../provider/PubSubTransportQueueFactory.java | 7 + .../RabbitMqTransportQueueFactory.java | 7 + .../ServiceBusTransportQueueFactory.java | 7 + .../provider/TbTransportQueueFactory.java | 3 + .../TbTransportQueueProducerProvider.java | 4 +- .../DefaultTransportRateLimitService.java | 27 ++-- .../service/DefaultTransportService.java | 15 +- .../DefaultNotificationSettingsService.java | 32 +++++ .../DefaultNotificationTargetService.java | 6 + .../notification/DefaultNotifications.java | 52 +++++++ .../notification/NotificationTargetDao.java | 3 + .../JpaNotificationTargetDao.java | 6 + .../util/AbstractBufferedRateExecutor.java | 2 +- .../util/limits/DefaultRateLimitService.java | 22 ++- .../server/dao/util/limits/LimitedApi.java | 67 --------- .../dao/util/limits/RateLimitService.java | 1 + .../src/main/resources/tb-vc-executor.yml | 7 + .../src/main/resources/tb-coap-transport.yml | 7 + .../src/main/resources/tb-http-transport.yml | 7 + .../src/main/resources/tb-lwm2m-transport.yml | 7 + .../src/main/resources/tb-mqtt-transport.yml | 7 + .../src/main/resources/tb-snmp-transport.yml | 7 + 49 files changed, 709 insertions(+), 107 deletions(-) create mode 100644 application/src/main/java/org/thingsboard/server/service/notification/rule/trigger/RateLimitsTriggerProcessor.java create mode 100644 common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApi.java create mode 100644 common/data/src/main/java/org/thingsboard/server/common/data/notification/info/RateLimitsNotificationInfo.java create mode 100644 common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/trigger/RateLimitsNotificationRuleTriggerConfig.java create mode 100644 common/message/src/main/java/org/thingsboard/server/common/msg/notification/trigger/RateLimitsTrigger.java delete mode 100644 dao/src/main/java/org/thingsboard/server/dao/util/limits/LimitedApi.java diff --git a/application/src/main/java/org/thingsboard/server/config/RateLimitProcessingFilter.java b/application/src/main/java/org/thingsboard/server/config/RateLimitProcessingFilter.java index 89f9b751fd..2ecc8590b8 100644 --- a/application/src/main/java/org/thingsboard/server/config/RateLimitProcessingFilter.java +++ b/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; } diff --git a/application/src/main/java/org/thingsboard/server/controller/AuthController.java b/application/src/main/java/org/thingsboard/server/controller/AuthController.java index 566b6900a4..4512334d61 100644 --- a/application/src/main/java/org/thingsboard/server/controller/AuthController.java +++ b/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; diff --git a/application/src/main/java/org/thingsboard/server/controller/plugin/TbWebSocketHandler.java b/application/src/main/java/org/thingsboard/server/controller/plugin/TbWebSocketHandler.java index bc51c5fdd1..91583136f6 100644 --- a/application/src/main/java/org/thingsboard/server/controller/plugin/TbWebSocketHandler.java +++ b/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; diff --git a/application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java b/application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java index 51d7435f0e..c6b9fe08a0 100644 --- a/application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java +++ b/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: diff --git a/application/src/main/java/org/thingsboard/server/service/install/DefaultSystemDataLoaderService.java b/application/src/main/java/org/thingsboard/server/service/install/DefaultSystemDataLoaderService.java index 497627f629..7965824274 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/DefaultSystemDataLoaderService.java +++ b/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); + } + } diff --git a/application/src/main/java/org/thingsboard/server/service/install/SystemDataLoaderService.java b/application/src/main/java/org/thingsboard/server/service/install/SystemDataLoaderService.java index fb6b28592c..08de12da78 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/SystemDataLoaderService.java +++ b/application/src/main/java/org/thingsboard/server/service/install/SystemDataLoaderService.java @@ -41,4 +41,6 @@ public interface SystemDataLoaderService { void createDefaultNotificationConfigs(); + void updateDefaultNotificationConfigs(); + } diff --git a/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationCenter.java b/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationCenter.java index 7215776170..cb77cc0948 100644 --- a/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationCenter.java +++ b/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; diff --git a/application/src/main/java/org/thingsboard/server/service/notification/rule/DefaultNotificationRuleProcessor.java b/application/src/main/java/org/thingsboard/server/service/notification/rule/DefaultNotificationRuleProcessor.java index 71a59026a2..3a5ae8339d 100644 --- a/application/src/main/java/org/thingsboard/server/service/notification/rule/DefaultNotificationRuleProcessor.java +++ b/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; diff --git a/application/src/main/java/org/thingsboard/server/service/notification/rule/trigger/RateLimitsTriggerProcessor.java b/application/src/main/java/org/thingsboard/server/service/notification/rule/trigger/RateLimitsTriggerProcessor.java new file mode 100644 index 0000000000..910ef4866f --- /dev/null +++ b/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 { + + 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; + } + +} diff --git a/application/src/main/java/org/thingsboard/server/service/security/auth/mfa/DefaultTwoFactorAuthService.java b/application/src/main/java/org/thingsboard/server/service/security/auth/mfa/DefaultTwoFactorAuthService.java index a58c801583..08717986bf 100644 --- a/application/src/main/java/org/thingsboard/server/service/security/auth/mfa/DefaultTwoFactorAuthService.java +++ b/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; diff --git a/application/src/main/java/org/thingsboard/server/service/sync/ie/DefaultEntitiesExportImportService.java b/application/src/main/java/org/thingsboard/server/service/sync/ie/DefaultEntitiesExportImportService.java index c8f8f945f9..3b37b0fe26 100644 --- a/application/src/main/java/org/thingsboard/server/service/sync/ie/DefaultEntitiesExportImportService.java +++ b/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; diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 392587c869..81589d2bce 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/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: 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 05ee66f327..c56adb5fd5 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 @@ -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); diff --git a/application/src/test/java/org/thingsboard/server/service/notification/MockNotificationSettingsService.java b/application/src/test/java/org/thingsboard/server/service/notification/MockNotificationSettingsService.java index a1cc82fd5a..a49f8dc8cb 100644 --- a/application/src/test/java/org/thingsboard/server/service/notification/MockNotificationSettingsService.java +++ b/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 diff --git a/application/src/test/java/org/thingsboard/server/service/notification/NotificationRuleApiTest.java b/application/src/test/java/org/thingsboard/server/service/notification/NotificationRuleApiTest.java index 6c4413d1e7..4155e2a4b3 100644 --- a/application/src/test/java/org/thingsboard/server/service/notification/NotificationRuleApiTest.java +++ b/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 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 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 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(); diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/notification/NotificationSettingsService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/notification/NotificationSettingsService.java index 46eb5ca54e..a5433915b3 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/notification/NotificationSettingsService.java +++ b/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); + } diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/notification/NotificationTargetService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/notification/NotificationTargetService.java index 05299c1675..bd442f6bb6 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/notification/NotificationTargetService.java +++ b/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 findNotificationTargetsByTenantIdAndIds(TenantId tenantId, List ids); + List findNotificationTargetsByTenantIdAndUsersFilterType(TenantId tenantId, UsersFilterType filterType); + PageData findRecipientsForNotificationTarget(TenantId tenantId, CustomerId customerId, NotificationTargetId targetId, PageLink pageLink); PageData findRecipientsForNotificationTargetConfig(TenantId tenantId, PlatformUsersNotificationTargetConfig targetConfig, PageLink pageLink); 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 new file mode 100644 index 0000000000..e709ac596e --- /dev/null +++ b/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 configExtractor; + @Getter + private final boolean perTenant; + @Getter + private boolean refillRateLimitIntervally; + @Getter + private String label; + + LimitedApi(Function 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); + } + +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/notification/NotificationType.java b/common/data/src/main/java/org/thingsboard/server/common/data/notification/NotificationType.java index de07d03c67..251fae2ac0 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/notification/NotificationType.java +++ b/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 } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/notification/info/RateLimitsNotificationInfo.java b/common/data/src/main/java/org/thingsboard/server/common/data/notification/info/RateLimitsNotificationInfo.java new file mode 100644 index 0000000000..57fc278cff --- /dev/null +++ b/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 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; + } + +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/trigger/NotificationRuleTriggerConfig.java b/common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/trigger/NotificationRuleTriggerConfig.java index b0eec28858..72028a616a 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/trigger/NotificationRuleTriggerConfig.java +++ b/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 { diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/trigger/NotificationRuleTriggerType.java b/common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/trigger/NotificationRuleTriggerType.java index dff86f4ba1..b591fe0245 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/trigger/NotificationRuleTriggerType.java +++ b/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; diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/trigger/RateLimitsNotificationRuleTriggerConfig.java b/common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/trigger/RateLimitsNotificationRuleTriggerConfig.java new file mode 100644 index 0000000000..f5ea83c47d --- /dev/null +++ b/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 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(",")); + } + +} diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/notification/trigger/RateLimitsTrigger.java b/common/message/src/main/java/org/thingsboard/server/common/msg/notification/trigger/RateLimitsTrigger.java new file mode 100644 index 0000000000..afb06bb8d1 --- /dev/null +++ b/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); + } + +} diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsTransportQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsTransportQueueFactory.java index c1ae9fe855..fa16ee9093 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsTransportQueueFactory.java +++ b/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> createTbCoreNotificationsMsgProducer() { + return new TbAwsSqsProducerTemplate<>(notificationAdmin, sqsSettings, coreSettings.getTopic()); + } + @Override public TbQueueConsumer> createTransportNotificationsConsumer() { return new TbAwsSqsConsumerTemplate<>(notificationAdmin, sqsSettings, transportNotificationSettings.getNotificationsTopic() + "_" + serviceInfoProvider.getServiceId(), diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/InMemoryTbTransportQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/InMemoryTbTransportQueueFactory.java index 60c464ecab..6b6253acc4 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/InMemoryTbTransportQueueFactory.java +++ b/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> createTbCoreNotificationsMsgProducer() { + return new InMemoryTbQueueProducer<>(storage, coreSettings.getTopic()); + } + @Override public TbQueueConsumer> createTransportNotificationsConsumer() { return new InMemoryTbQueueConsumer<>(storage, transportNotificationSettings.getNotificationsTopic() + "." + serviceInfoProvider.getServiceId()); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbTransportQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbTransportQueueFactory.java index 51e3dc99c7..300e0a44b6 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbTransportQueueFactory.java +++ b/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> createTbCoreNotificationsMsgProducer() { + TbKafkaProducerTemplate.TbKafkaProducerTemplateBuilder> 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> createTransportNotificationsConsumer() { TbKafkaConsumerTemplate.TbKafkaConsumerTemplateBuilder> responseBuilder = TbKafkaConsumerTemplate.builder(); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTransportQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTransportQueueFactory.java index d5665e2fbe..5cafaa4c48 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTransportQueueFactory.java +++ b/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> createTbCoreNotificationsMsgProducer() { + return new TbPubSubProducerTemplate<>(notificationAdmin, pubSubSettings, coreSettings.getTopic()); + } + @Override public TbQueueConsumer> createTransportNotificationsConsumer() { return new TbPubSubConsumerTemplate<>(notificationAdmin, pubSubSettings, diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/RabbitMqTransportQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/RabbitMqTransportQueueFactory.java index f11a877fc5..bcbd60a14e 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/RabbitMqTransportQueueFactory.java +++ b/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> createTbCoreNotificationsMsgProducer() { + return new TbRabbitMqProducerTemplate<>(notificationAdmin, rabbitMqSettings, coreSettings.getTopic()); + } + @Override public TbQueueConsumer> createTransportNotificationsConsumer() { return new TbRabbitMqConsumerTemplate<>(notificationAdmin, rabbitMqSettings, transportNotificationSettings.getNotificationsTopic() + "." + serviceInfoProvider.getServiceId(), diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/ServiceBusTransportQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/ServiceBusTransportQueueFactory.java index 15e661f5b0..0a4bb59663 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/ServiceBusTransportQueueFactory.java +++ b/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> createTbCoreNotificationsMsgProducer() { + return new TbServiceBusProducerTemplate<>(notificationAdmin, serviceBusSettings, coreSettings.getTopic()); + } + @Override public TbQueueConsumer> createTransportNotificationsConsumer() { return new TbServiceBusConsumerTemplate<>(notificationAdmin, serviceBusSettings, diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbTransportQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbTransportQueueFactory.java index feff3830fc..c7d2f78838 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbTransportQueueFactory.java +++ b/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> createTbCoreMsgProducer(); + TbQueueProducer> createTbCoreNotificationsMsgProducer(); + TbQueueConsumer> createTransportNotificationsConsumer(); } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbTransportQueueProducerProvider.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbTransportQueueProducerProvider.java index e9177d8b1c..afbcd718a5 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbTransportQueueProducerProvider.java +++ b/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> toRuleEngine; private TbQueueProducer> toTbCore; + private TbQueueProducer> toTbCoreNotifications; private TbQueueProducer> 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> getTbCoreNotificationsMsgProducer() { - throw new RuntimeException("Not Implemented! Should not be used by Transport!"); + return toTbCoreNotifications; } @Override diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/DefaultTransportRateLimitService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/DefaultTransportRateLimitService.java index 6539f2d22c..47f7717212 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/DefaultTransportRateLimitService.java +++ b/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; + }); } } diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java index b495448d90..69e4de7b74 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java +++ b/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> transportApiRequestTemplate; protected TbQueueProducer> 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; } } diff --git a/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationSettingsService.java b/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationSettingsService.java index 0e25a0cab3..4262decfed 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationSettingsService.java +++ b/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); diff --git a/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationTargetService.java b/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationTargetService.java index 78d5a4b454..97a7ee0e27 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationTargetService.java +++ b/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 findNotificationTargetsByTenantIdAndUsersFilterType(TenantId tenantId, UsersFilterType filterType) { + return notificationTargetDao.findByTenantIdAndUsersFilterType(tenantId, filterType); + } + @Override public PageData findRecipientsForNotificationTarget(TenantId tenantId, CustomerId customerId, NotificationTargetId targetId, PageLink pageLink) { NotificationTarget notificationTarget = findNotificationTargetById(tenantId, targetId); diff --git a/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotifications.java b/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotifications.java index 01f9f68ae4..c8f458b40d 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotifications.java +++ b/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) diff --git a/dao/src/main/java/org/thingsboard/server/dao/notification/NotificationTargetDao.java b/dao/src/main/java/org/thingsboard/server/dao/notification/NotificationTargetDao.java index 912e873bcc..3e542cc50f 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/notification/NotificationTargetDao.java +++ b/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, TenantEn List findByTenantIdAndIds(TenantId tenantId, List ids); + List findByTenantIdAndUsersFilterType(TenantId tenantId, UsersFilterType filterType); + void removeByTenantId(TenantId tenantId); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/notification/JpaNotificationTargetDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/notification/JpaNotificationTargetDao.java index d976802fd2..64c679c845 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/notification/JpaNotificationTargetDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/notification/JpaNotificationTargetDao.java @@ -66,6 +66,12 @@ public class JpaNotificationTargetDao extends JpaAbstractDao 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()); 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 64b4eb3666..4c8d0c5afe 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 @@ -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; diff --git a/dao/src/main/java/org/thingsboard/server/dao/util/limits/DefaultRateLimitService.java b/dao/src/main/java/org/thingsboard/server/dao/util/limits/DefaultRateLimitService.java index 918f16ef19..ede6855d5e 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/util/limits/DefaultRateLimitService.java +++ b/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 diff --git a/dao/src/main/java/org/thingsboard/server/dao/util/limits/LimitedApi.java b/dao/src/main/java/org/thingsboard/server/dao/util/limits/LimitedApi.java deleted file mode 100644 index ee79230d8a..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/util/limits/LimitedApi.java +++ /dev/null @@ -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 configExtractor; - @Getter - private final boolean refillRateLimitIntervally; - - LimitedApi(Function configExtractor) { - this((profileConfiguration, level) -> configExtractor.apply(profileConfiguration)); - } - - LimitedApi(BiFunction 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"); - } - } - -} diff --git a/dao/src/main/java/org/thingsboard/server/dao/util/limits/RateLimitService.java b/dao/src/main/java/org/thingsboard/server/dao/util/limits/RateLimitService.java index c3f2fd179f..1f6e87111c 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/util/limits/RateLimitService.java +++ b/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 { diff --git a/msa/vc-executor/src/main/resources/tb-vc-executor.yml b/msa/vc-executor/src/main/resources/tb-vc-executor.yml index 9d5ad35388..75f9e09d3c 100644 --- a/msa/vc-executor/src/main/resources/tb-vc-executor.yml +++ b/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}" diff --git a/transport/coap/src/main/resources/tb-coap-transport.yml b/transport/coap/src/main/resources/tb-coap-transport.yml index 7ea553fe5c..8fe079859a 100644 --- a/transport/coap/src/main/resources/tb-coap-transport.yml +++ b/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}" diff --git a/transport/http/src/main/resources/tb-http-transport.yml b/transport/http/src/main/resources/tb-http-transport.yml index 346ec48eae..f05db08643 100644 --- a/transport/http/src/main/resources/tb-http-transport.yml +++ b/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}" diff --git a/transport/lwm2m/src/main/resources/tb-lwm2m-transport.yml b/transport/lwm2m/src/main/resources/tb-lwm2m-transport.yml index 4e8167d89d..745a7d126a 100644 --- a/transport/lwm2m/src/main/resources/tb-lwm2m-transport.yml +++ b/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}" diff --git a/transport/mqtt/src/main/resources/tb-mqtt-transport.yml b/transport/mqtt/src/main/resources/tb-mqtt-transport.yml index 1e0b1ebcd4..3795c533a7 100644 --- a/transport/mqtt/src/main/resources/tb-mqtt-transport.yml +++ b/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}" diff --git a/transport/snmp/src/main/resources/tb-snmp-transport.yml b/transport/snmp/src/main/resources/tb-snmp-transport.yml index 9f086bcbc5..11dcc96010 100644 --- a/transport/snmp/src/main/resources/tb-snmp-transport.yml +++ b/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}"