From 409976c6a316755d02af611e41599b05180b8026 Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Fri, 23 Dec 2022 18:47:40 +0200 Subject: [PATCH] Multiple notification targets for request; improvements and refactoring --- .../main/data/upgrade/3.4.3/schema_update.sql | 14 +- .../server/actors/ActorSystemContext.java | 5 + .../actors/ruleChain/DefaultTbContext.java | 5 + .../server/controller/BaseController.java | 12 +- .../controller/NotificationController.java | 3 +- .../NotificationRuleController.java | 7 +- .../NotificationTargetController.java | 20 +-- .../NotificationTemplateController.java | 3 +- .../NotificationExecutorService.java | 29 +++ .../DefaultSystemDataLoaderService.java | 5 + .../install/SystemDataLoaderService.java | 3 + .../DefaultNotificationManager.java | 80 +++++---- ...aultNotificationRuleProcessingService.java | 53 ++++-- .../DefaultNotificationSchedulerService.java | 59 ++++-- .../NotificationProcessingContext.java | 34 ++-- .../NotificationRuleProcessingService.java | 3 - .../NotificationSchedulerService.java | 2 - .../channels/SlackNotificationChannel.java | 2 +- .../queue/DefaultTbCoreConsumerService.java | 15 +- .../DefaultTbRuleEngineConsumerService.java | 6 +- .../processing/AbstractConsumerService.java | 7 +- .../service/slack/DefaultSlackService.java | 2 +- .../ttl/NotificationsCleanUpService.java | 20 ++- .../server/controller/AbstractWebTest.java | 27 +-- .../notification/NotificationApiTest.java | 23 +-- .../NotificationTargetApiTest.java | 168 ++++++++++++++++++ .../NotificationTemplateApiTest.java | 90 ++++++++++ .../notification/NotificationsClient.java | 86 --------- common/cluster-api/src/main/proto/queue.proto | 1 - .../NotificationRequestService.java | 3 + .../notification/NotificationRuleService.java | 4 +- .../NotificationTargetService.java | 7 +- .../notification/NotificationRequest.java | 14 +- .../NotificationRequestConfig.java | 4 +- .../NotificationRequestStats.java | 15 +- .../NonConfirmedNotificationEscalation.java | 5 +- ...CustomerUsersNotificationTargetConfig.java | 16 +- .../targets/NotificationTarget.java | 3 +- .../DeliveryMethodNotificationTemplate.java | 3 + ...ailDeliveryMethodNotificationTemplate.java | 3 + .../template/NotificationTemplate.java | 1 + .../template/NotificationTemplateConfig.java | 19 +- .../template}/SlackConversation.java | 2 +- ...ackDeliveryMethodNotificationTemplate.java | 2 + .../server/dao/model/BaseSqlEntity.java | 44 ++++- .../server/dao/model/ModelConstants.java | 2 +- .../dao/model/sql/AbstractAlarmEntity.java | 2 +- .../dao/model/sql/NotificationEntity.java | 5 +- .../model/sql/NotificationRequestEntity.java | 27 ++- .../dao/model/sql/NotificationRuleEntity.java | 17 +- .../model/sql/NotificationTargetEntity.java | 4 +- .../model/sql/NotificationTemplateEntity.java | 4 +- .../DefaultNotificationRequestService.java | 9 +- .../DefaultNotificationRuleService.java | 8 +- .../DefaultNotificationTargetService.java | 38 ++-- .../DefaultNotificationTemplateService.java | 19 +- .../notification/NotificationRequestDao.java | 10 ++ .../dao/notification/NotificationRuleDao.java | 3 + .../dao/service/ConstraintValidator.java | 27 ++- .../server/dao/service/NoXssValidator.java | 1 + .../JpaNotificationRequestDao.java | 24 +++ .../notification/JpaNotificationRuleDao.java | 6 + .../NotificationRequestRepository.java | 10 ++ .../NotificationRuleRepository.java | 2 + .../insert/sql/SqlPartitioningRepository.java | 5 +- .../main/resources/sql/schema-entities.sql | 14 +- .../rule/engine/api/TbContext.java | 2 + .../rule/engine/api/slack/SlackService.java | 1 + .../notification/TbNotificationNode.java | 6 +- .../TbNotificationNodeConfiguration.java | 7 +- .../rule/engine/notification/TbSlackNode.java | 2 +- .../TbSlackNodeConfiguration.java | 2 +- 72 files changed, 844 insertions(+), 342 deletions(-) create mode 100644 application/src/main/java/org/thingsboard/server/service/executors/NotificationExecutorService.java create mode 100644 application/src/test/java/org/thingsboard/server/service/notification/NotificationTargetApiTest.java create mode 100644 application/src/test/java/org/thingsboard/server/service/notification/NotificationTemplateApiTest.java delete mode 100644 application/src/test/java/org/thingsboard/server/service/notification/NotificationsClient.java rename {rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/slack => common/data/src/main/java/org/thingsboard/server/common/data/notification/template}/SlackConversation.java (92%) diff --git a/application/src/main/data/upgrade/3.4.3/schema_update.sql b/application/src/main/data/upgrade/3.4.3/schema_update.sql index b9aaaa5c7b..fa113d6c5d 100644 --- a/application/src/main/data/upgrade/3.4.3/schema_update.sql +++ b/application/src/main/data/upgrade/3.4.3/schema_update.sql @@ -17,7 +17,7 @@ CREATE TABLE IF NOT EXISTS notification_target ( id UUID NOT NULL CONSTRAINT notification_target_pkey PRIMARY KEY, created_time BIGINT NOT NULL, - tenant_id UUID NOT NULL CONSTRAINT fk_notification_target_tenant_id REFERENCES tenant(id) ON DELETE CASCADE, + tenant_id UUID NULL CONSTRAINT fk_notification_target_tenant_id REFERENCES tenant(id) ON DELETE CASCADE, name VARCHAR(255) NOT NULL, configuration VARCHAR(10000) NOT NULL ); @@ -26,7 +26,7 @@ CREATE INDEX IF NOT EXISTS idx_notification_target_tenant_id_created_time ON not CREATE TABLE IF NOT EXISTS notification_template ( id UUID NOT NULL CONSTRAINT notification_template_pkey PRIMARY KEY, created_time BIGINT NOT NULL, - tenant_id UUID NOT NULL CONSTRAINT fk_notification_template_tenant_id REFERENCES tenant(id) ON DELETE CASCADE, + tenant_id UUID NULL CONSTRAINT fk_notification_template_tenant_id REFERENCES tenant(id) ON DELETE CASCADE, name VARCHAR(255) NOT NULL, notification_type VARCHAR(255) NOT NULL, configuration VARCHAR(10000) NOT NULL @@ -35,7 +35,7 @@ CREATE TABLE IF NOT EXISTS notification_template ( CREATE TABLE IF NOT EXISTS notification_rule ( id UUID NOT NULL CONSTRAINT notification_rule_pkey PRIMARY KEY, created_time BIGINT NOT NULL, - tenant_id UUID NOT NULL CONSTRAINT fk_notification_rule_tenant_id REFERENCES tenant(id) ON DELETE CASCADE, + tenant_id UUID NULL CONSTRAINT fk_notification_rule_tenant_id REFERENCES tenant(id) ON DELETE CASCADE, name VARCHAR(255) NOT NULL, template_id UUID NOT NULL CONSTRAINT fk_notification_rule_template_id REFERENCES notification_template(id), delivery_methods VARCHAR(255) NOT NULL, @@ -45,16 +45,16 @@ CREATE TABLE IF NOT EXISTS notification_rule ( CREATE TABLE IF NOT EXISTS notification_request ( id UUID NOT NULL CONSTRAINT notification_request_pkey PRIMARY KEY, created_time BIGINT NOT NULL, - tenant_id UUID NOT NULL CONSTRAINT fk_notification_request_tenant_id REFERENCES tenant(id) ON DELETE CASCADE, - target_id UUID NOT NULL CONSTRAINT fk_notification_request_target_id REFERENCES notification_target(id), - template_id UUID NOT NULL CONSTRAINT fk_notification_request_template_id REFERENCES notification_template(id), + tenant_id UUID NULL CONSTRAINT fk_notification_request_tenant_id REFERENCES tenant(id) ON DELETE CASCADE, + targets VARCHAR(255) NOT NULL, + template_id UUID NOT NULL, info VARCHAR(1000), delivery_methods VARCHAR(255), additional_config VARCHAR(1000), originator_type VARCHAR(32) NOT NULL, originator_entity_id UUID, originator_entity_type VARCHAR(32), - rule_id UUID NULL CONSTRAINT fk_notification_request_rule_id REFERENCES notification_rule(id), + rule_id UUID NULL, status VARCHAR(32), stats VARCHAR(1000) ); diff --git a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java index 9035a77ea3..8e754bf1b1 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java +++ b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java @@ -91,6 +91,7 @@ import org.thingsboard.server.service.edge.rpc.EdgeRpcService; import org.thingsboard.server.service.entitiy.entityview.TbEntityViewService; import org.thingsboard.server.service.executors.DbCallbackExecutorService; import org.thingsboard.server.service.executors.ExternalCallExecutorService; +import org.thingsboard.server.service.executors.NotificationExecutorService; import org.thingsboard.server.service.executors.SharedEventLoopGroupService; import org.thingsboard.server.service.mail.MailExecutorService; import org.thingsboard.server.service.profile.TbAssetProfileCache; @@ -304,6 +305,10 @@ public class ActorSystemContext { @Getter private ExternalCallExecutorService externalCallExecutorService; + @Autowired + @Getter + private NotificationExecutorService notificationExecutor; + @Autowired @Getter private SharedEventLoopGroupService sharedEventLoopGroupService; diff --git a/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java b/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java index eeb5a8fe1b..6c70eb12b9 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java +++ b/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java @@ -472,6 +472,11 @@ class DefaultTbContext implements TbContext { return mainCtx.getExternalCallExecutorService(); } + @Override + public ListeningExecutor getNotificationExecutor() { + return mainCtx.getNotificationExecutor(); + } + @Override @Deprecated public ScriptEngine createJsScriptEngine(String script, String... argNames) { diff --git a/application/src/main/java/org/thingsboard/server/controller/BaseController.java b/application/src/main/java/org/thingsboard/server/controller/BaseController.java index 841aff37e5..e31c690a60 100644 --- a/application/src/main/java/org/thingsboard/server/controller/BaseController.java +++ b/application/src/main/java/org/thingsboard/server/controller/BaseController.java @@ -377,11 +377,17 @@ public abstract class BaseController { * */ @ExceptionHandler(MethodArgumentNotValidException.class) public void handleValidationError(MethodArgumentNotValidException e, HttpServletResponse response) { - String errorMessage = "Validation error: " + e.getBindingResult().getAllErrors().stream() - .map(DefaultMessageSourceResolvable::getDefaultMessage) + String errorMessage = "Validation error: " + e.getFieldErrors().stream() + .map(fieldError -> { + String property = fieldError.getField(); + if (property.equals("valid") || StringUtils.endsWith(property, ".valid")) { // when custom @AssertTrue is used + property = ""; + } + return (!property.isEmpty() ? (property + " ") : "") + fieldError.getDefaultMessage(); + }) .collect(Collectors.joining(", ")); ThingsboardException thingsboardException = new ThingsboardException(errorMessage, ThingsboardErrorCode.BAD_REQUEST_PARAMS); - handleThingsboardException(thingsboardException, response); + handleControllerException(thingsboardException, response); } T checkNotNull(T reference) throws ThingsboardException { diff --git a/application/src/main/java/org/thingsboard/server/controller/NotificationController.java b/application/src/main/java/org/thingsboard/server/controller/NotificationController.java index 3c7e8865b9..3d6153afa0 100644 --- a/application/src/main/java/org/thingsboard/server/controller/NotificationController.java +++ b/application/src/main/java/org/thingsboard/server/controller/NotificationController.java @@ -29,7 +29,7 @@ import org.springframework.web.bind.annotation.RequestMapping; import org.springframework.web.bind.annotation.RequestParam; import org.springframework.web.bind.annotation.RestController; import org.thingsboard.rule.engine.api.NotificationManager; -import org.thingsboard.rule.engine.api.slack.SlackConversation; +import org.thingsboard.server.common.data.notification.template.SlackConversation; import org.thingsboard.rule.engine.api.slack.SlackService; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.exception.ThingsboardException; @@ -102,6 +102,7 @@ public class NotificationController extends BaseController { notificationRequest.setOriginatorType(NotificationOriginatorType.ADMIN); notificationRequest.setOriginatorEntityId(user.getId()); + notificationRequest.setOriginatorEntity(user); if (notificationRequest.getInfo() != null && notificationRequest.getInfo().getOriginatorType() != null) { throw new IllegalArgumentException("Unsupported notification info type"); } diff --git a/application/src/main/java/org/thingsboard/server/controller/NotificationRuleController.java b/application/src/main/java/org/thingsboard/server/controller/NotificationRuleController.java index 0ed5e1bfaf..8d5a4bde34 100644 --- a/application/src/main/java/org/thingsboard/server/controller/NotificationRuleController.java +++ b/application/src/main/java/org/thingsboard/server/controller/NotificationRuleController.java @@ -33,6 +33,7 @@ import org.thingsboard.server.common.data.id.NotificationRuleId; import org.thingsboard.server.common.data.notification.rule.NotificationRule; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageLink; +import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; import org.thingsboard.server.dao.notification.NotificationRuleService; import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.security.model.SecurityUser; @@ -78,10 +79,12 @@ public class NotificationRuleController extends BaseController { @DeleteMapping("/rule/{id}") @PreAuthorize("hasAnyAuthority('TENANT_ADMIN')") - public void deleteNotificationRule(@PathVariable UUID id) throws Exception { + public void deleteNotificationRule(@PathVariable UUID id, + @AuthenticationPrincipal SecurityUser user) throws Exception { NotificationRuleId notificationRuleId = new NotificationRuleId(id); NotificationRule notificationRule = checkEntityId(notificationRuleId, notificationRuleService::findNotificationRuleById, Operation.DELETE); - doDeleteAndLog(EntityType.NOTIFICATION_RULE, notificationRule, notificationRuleService::deleteNotificationRule); + doDeleteAndLog(EntityType.NOTIFICATION_RULE, notificationRule, notificationRuleService::deleteNotificationRuleById); + tbClusterService.broadcastEntityStateChangeEvent(user.getTenantId(), notificationRuleId, ComponentLifecycleEvent.DELETED); } } diff --git a/application/src/main/java/org/thingsboard/server/controller/NotificationTargetController.java b/application/src/main/java/org/thingsboard/server/controller/NotificationTargetController.java index ee984d9a9a..00ac6f4a93 100644 --- a/application/src/main/java/org/thingsboard/server/controller/NotificationTargetController.java +++ b/application/src/main/java/org/thingsboard/server/controller/NotificationTargetController.java @@ -33,8 +33,8 @@ import org.thingsboard.server.common.data.exception.ThingsboardException; import org.thingsboard.server.common.data.id.NotificationTargetId; import org.thingsboard.server.common.data.notification.targets.NotificationTarget; import org.thingsboard.server.common.data.notification.targets.NotificationTargetConfig; -import org.thingsboard.server.common.data.notification.targets.NotificationTargetConfigType; import org.thingsboard.server.common.data.page.PageData; +import org.thingsboard.server.common.data.page.PageDataIterable; import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.dao.notification.NotificationTargetService; import org.thingsboard.server.queue.util.TbCoreComponent; @@ -42,6 +42,7 @@ import org.thingsboard.server.service.security.model.SecurityUser; import org.thingsboard.server.service.security.permission.Operation; import org.thingsboard.server.service.security.permission.Resource; +import javax.validation.Valid; import java.util.UUID; @RestController @@ -55,17 +56,16 @@ public class NotificationTargetController extends BaseController { @PostMapping("/target") @PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN')") - public NotificationTarget saveNotificationTarget(@RequestBody NotificationTarget notificationTarget, + public NotificationTarget saveNotificationTarget(@RequestBody @Valid NotificationTarget notificationTarget, @AuthenticationPrincipal SecurityUser user) throws Exception { checkEntity(notificationTarget.getId(), notificationTarget, Resource.NOTIFICATION_TARGET); if (!user.isSystemAdmin()) { NotificationTargetConfig targetConfig = notificationTarget.getConfiguration(); - if (targetConfig.getType() == NotificationTargetConfigType.SINGLE_USER || - targetConfig.getType() == NotificationTargetConfigType.USER_LIST) { - PageData recipients = notificationTargetService.findRecipientsForNotificationTargetConfig(user.getTenantId(), notificationTarget.getConfiguration(), null); - for (User recipient : recipients.getData()) { - accessControlService.checkPermission(user, Resource.USER, Operation.READ, recipient.getId(), recipient); - } + PageDataIterable recipients = new PageDataIterable<>(pageLink -> { + return notificationTargetService.findRecipientsForNotificationTargetConfig(user.getTenantId(), null, targetConfig, pageLink); + }, 200); + for (User recipient : recipients) { + accessControlService.checkPermission(user, Resource.USER, Operation.READ, recipient.getId(), recipient); } } @@ -86,7 +86,7 @@ public class NotificationTargetController extends BaseController { @RequestParam int page, @AuthenticationPrincipal SecurityUser user) throws ThingsboardException { PageLink pageLink = createPageLink(pageSize, page, null, null, null); - PageData recipients = notificationTargetService.findRecipientsForNotificationTargetConfig(user.getTenantId(), notificationTarget.getConfiguration(), pageLink); + PageData recipients = notificationTargetService.findRecipientsForNotificationTargetConfig(user.getTenantId(), null, notificationTarget.getConfiguration(), pageLink); if (!user.isSystemAdmin()) { for (User recipient : recipients.getData()) { accessControlService.checkPermission(user, Resource.USER, Operation.READ, recipient.getId(), recipient); @@ -112,7 +112,7 @@ public class NotificationTargetController extends BaseController { public void deleteNotificationTarget(@PathVariable UUID id) throws Exception { NotificationTargetId notificationTargetId = new NotificationTargetId(id); NotificationTarget notificationTarget = checkEntityId(notificationTargetId, notificationTargetService::findNotificationTargetById, Operation.DELETE); - doDeleteAndLog(EntityType.NOTIFICATION_TARGET, notificationTarget, notificationTargetService::deleteNotificationTarget); + doDeleteAndLog(EntityType.NOTIFICATION_TARGET, notificationTarget, notificationTargetService::deleteNotificationTargetById); } } diff --git a/application/src/main/java/org/thingsboard/server/controller/NotificationTemplateController.java b/application/src/main/java/org/thingsboard/server/controller/NotificationTemplateController.java index 3d5e0c8332..5194d57f84 100644 --- a/application/src/main/java/org/thingsboard/server/controller/NotificationTemplateController.java +++ b/application/src/main/java/org/thingsboard/server/controller/NotificationTemplateController.java @@ -33,6 +33,7 @@ import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.security.permission.Operation; import org.thingsboard.server.service.security.permission.Resource; +import javax.validation.Valid; import java.util.UUID; @RestController @@ -45,7 +46,7 @@ public class NotificationTemplateController extends BaseController { @PostMapping("/template") @PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN')") - public NotificationTemplate saveNotificationTemplate(@RequestBody NotificationTemplate notificationTemplate) throws Exception { + public NotificationTemplate saveNotificationTemplate(@RequestBody @Valid NotificationTemplate notificationTemplate) throws Exception { checkEntity(notificationTemplate.getId(), notificationTemplate, Resource.NOTIFICATION_TEMPLATE); return doSaveAndLog(EntityType.NOTIFICATION_TEMPLATE, notificationTemplate, notificationTemplateService::saveNotificationTemplate); } diff --git a/application/src/main/java/org/thingsboard/server/service/executors/NotificationExecutorService.java b/application/src/main/java/org/thingsboard/server/service/executors/NotificationExecutorService.java new file mode 100644 index 0000000000..15f1a27184 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/executors/NotificationExecutorService.java @@ -0,0 +1,29 @@ +/** + * Copyright © 2016-2022 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.executors; + +import org.springframework.stereotype.Component; +import org.thingsboard.common.util.AbstractListeningExecutor; + +@Component +public class NotificationExecutorService extends AbstractListeningExecutor { + + @Override + protected int getThreadPollSize() { + return 10; // FIXME [viacheslav] + } + +} 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 9507912c39..dd1f153f51 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 @@ -671,4 +671,9 @@ public class DefaultSystemDataLoaderService implements SystemDataLoaderService { } } + @Override + public void createNotificationConfigs() { + // create default notification targets: Alarm's customer + } + } 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 041351f5e4..b5a54abed0 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 @@ -38,4 +38,7 @@ public interface SystemDataLoaderService { void deleteSystemWidgetBundle(String bundleAlias) throws Exception; void createQueues(); + + void createNotificationConfigs(); + } diff --git a/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationManager.java b/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationManager.java index f04cced9f6..28421bc429 100644 --- a/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationManager.java +++ b/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationManager.java @@ -27,8 +27,10 @@ import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.id.NotificationId; import org.thingsboard.server.common.data.id.NotificationRequestId; +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.notification.AlreadySentException; import org.thingsboard.server.common.data.notification.Notification; import org.thingsboard.server.common.data.notification.NotificationDeliveryMethod; import org.thingsboard.server.common.data.notification.NotificationRequest; @@ -37,7 +39,7 @@ import org.thingsboard.server.common.data.notification.NotificationRequestStatus import org.thingsboard.server.common.data.notification.NotificationStatus; import org.thingsboard.server.common.data.notification.settings.NotificationSettings; import org.thingsboard.server.common.data.notification.template.DeliveryMethodNotificationTemplate; -import org.thingsboard.server.common.data.page.PageLink; +import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; @@ -85,7 +87,7 @@ public class DefaultNotificationManager extends AbstractSubscriptionService impl @Override public NotificationRequest processNotificationRequest(TenantId tenantId, NotificationRequest notificationRequest) { - log.debug("Processing notification request (tenant id: {}, notification target id: {})", tenantId, notificationRequest.getTargetId()); + log.debug("Processing notification request (tenant id: {}, notification targets: {})", tenantId, notificationRequest.getTargets()); notificationRequest.setTenantId(tenantId); NotificationSettings settings = notificationSettingsService.findNotificationSettings(tenantId); notificationRequest.getDeliveryMethods().forEach(deliveryMethod -> { @@ -99,62 +101,63 @@ public class DefaultNotificationManager extends AbstractSubscriptionService impl if (config.getSendingDelayInSec() > 0 && notificationRequest.getId() == null) { notificationRequest.setStatus(NotificationRequestStatus.SCHEDULED); NotificationRequest savedNotificationRequest = notificationRequestService.saveNotificationRequest(tenantId, notificationRequest); - forwardToNotificationSchedulerService(tenantId, savedNotificationRequest.getId(), false); + forwardToNotificationSchedulerService(tenantId, savedNotificationRequest.getId()); return savedNotificationRequest; } } - if (notificationTargetService.findRecipientsForNotificationTarget(tenantId, notificationRequest.getTargetId(), new PageLink(1)) - .getTotalElements() == 0) { - throw new IllegalArgumentException("No target recipients"); - } notificationRequest.setStatus(NotificationRequestStatus.PROCESSED); NotificationRequest savedNotificationRequest = notificationRequestService.saveNotificationRequest(tenantId, notificationRequest); NotificationProcessingContext ctx = NotificationProcessingContext.builder() .tenantId(tenantId) - .settings(settings) .request(savedNotificationRequest) - .additionalTemplateContext(notificationRequest.getTemplateContext()) + .settings(settings) .build(); ctx.init(notificationTemplateService); - DaoUtil.processBatches(pageLink -> { - return notificationTargetService.findRecipientsForNotificationTarget(tenantId, notificationRequest.getTargetId(), pageLink); - }, 200, recipientsBatch -> { - List> results = new ArrayList<>(); - for (NotificationDeliveryMethod deliveryMethod : savedNotificationRequest.getDeliveryMethods()) { - NotificationChannel notificationChannel = channels.get(deliveryMethod); - log.debug("Sending {} notifications for request {} to recipients batch", deliveryMethod, savedNotificationRequest.getId()); + for (NotificationTargetId targetId : notificationRequest.getTargets()) { + DaoUtil.processBatches(pageLink -> { + return notificationTargetService.findRecipientsForNotificationTarget(tenantId, ctx.getOriginatorCustomerId(), targetId, pageLink); + }, 200, recipientsBatch -> { + List> results = new ArrayList<>(); + for (NotificationDeliveryMethod deliveryMethod : savedNotificationRequest.getDeliveryMethods()) { + NotificationChannel notificationChannel = channels.get(deliveryMethod); + log.debug("Sending {} notifications for request {} to recipients batch", deliveryMethod, savedNotificationRequest.getId()); - List recipients = recipientsBatch.getData(); - for (User recipient : recipients) { - ListenableFuture resultFuture = processForRecipient(notificationChannel, recipient, ctx); - DonAsynchron.withCallback(resultFuture, result -> { - ctx.getStats().reportSent(deliveryMethod); - }, error -> { - ctx.getStats().reportError(deliveryMethod, recipient, error); - }, dbCallbackExecutorService); - results.add(resultFuture); + List recipients = recipientsBatch.getData(); + for (User recipient : recipients) { + ListenableFuture resultFuture = processForRecipient(notificationChannel, recipient, ctx); + DonAsynchron.withCallback(resultFuture, result -> { + ctx.getStats().reportSent(deliveryMethod, recipient); + }, error -> { + ctx.getStats().reportError(deliveryMethod, recipient, error); + }); + results.add(resultFuture); + } } - } - Futures.allAsList(results).addListener(() -> { - try { - notificationRequestService.updateNotificationRequestStats(tenantId, savedNotificationRequest.getId(), ctx.getStats()); - } catch (Exception e) { - log.error("Failed to update stats for notification request {}", savedNotificationRequest.getId(), e); - } - }, dbCallbackExecutorService); - }); + Futures.allAsList(results).addListener(() -> { + try { + notificationRequestService.updateNotificationRequestStats(tenantId, savedNotificationRequest.getId(), ctx.getStats()); + } catch (Exception e) { + log.error("Failed to update stats for notification request {}", savedNotificationRequest.getId(), e); + } + }, dbCallbackExecutorService); + }); + } return savedNotificationRequest; } private ListenableFuture processForRecipient(NotificationChannel notificationChannel, User recipient, NotificationProcessingContext ctx) { + NotificationDeliveryMethod deliveryMethod = notificationChannel.getDeliveryMethod(); + if (ctx.getStats().contains(deliveryMethod, recipient.getId())) { + return Futures.immediateFailedFuture(new AlreadySentException()); + } String text; try { - DeliveryMethodNotificationTemplate template = ctx.getTemplate(notificationChannel.getDeliveryMethod()); + DeliveryMethodNotificationTemplate template = ctx.getTemplate(deliveryMethod); text = TbNodeUtils.processTemplate(template.getBody(), ctx.createTemplateContext(recipient)); } catch (Exception e) { return Futures.immediateFailedFuture(e); @@ -162,14 +165,13 @@ public class DefaultNotificationManager extends AbstractSubscriptionService impl return notificationChannel.sendNotification(recipient, text, ctx); } - private void forwardToNotificationSchedulerService(TenantId tenantId, NotificationRequestId notificationRequestId, boolean deleted) { + private void forwardToNotificationSchedulerService(TenantId tenantId, NotificationRequestId notificationRequestId) { TransportProtos.NotificationSchedulerServiceMsg.Builder msg = TransportProtos.NotificationSchedulerServiceMsg.newBuilder() .setTenantIdMSB(tenantId.getId().getMostSignificantBits()) .setTenantIdLSB(tenantId.getId().getLeastSignificantBits()) .setRequestIdMSB(notificationRequestId.getId().getMostSignificantBits()) .setRequestIdLSB(notificationRequestId.getId().getLeastSignificantBits()) - .setTs(System.currentTimeMillis()) - .setDeleted(deleted); + .setTs(System.currentTimeMillis()); TransportProtos.ToCoreMsg toCoreMsg = TransportProtos.ToCoreMsg.newBuilder() .setNotificationSchedulerServiceMsg(msg) .build(); @@ -216,7 +218,7 @@ public class DefaultNotificationManager extends AbstractSubscriptionService impl .notificationRequestId(notificationRequestId) .deleted(true) .build()); - forwardToNotificationSchedulerService(tenantId, notificationRequestId, true); + clusterService.broadcastEntityStateChangeEvent(tenantId, notificationRequestId, ComponentLifecycleEvent.DELETED); } @Override diff --git a/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationRuleProcessingService.java b/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationRuleProcessingService.java index 5f97d47093..e6ef9bd8ba 100644 --- a/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationRuleProcessingService.java +++ b/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationRuleProcessingService.java @@ -21,9 +21,12 @@ import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.annotation.Lazy; +import org.springframework.context.event.EventListener; import org.springframework.stereotype.Service; import org.thingsboard.rule.engine.api.NotificationManager; +import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.alarm.Alarm; +import org.thingsboard.server.common.data.id.NotificationRequestId; import org.thingsboard.server.common.data.id.NotificationRuleId; import org.thingsboard.server.common.data.id.NotificationTargetId; import org.thingsboard.server.common.data.id.TenantId; @@ -32,13 +35,16 @@ import org.thingsboard.server.common.data.notification.NotificationInfo; import org.thingsboard.server.common.data.notification.NotificationOriginatorType; import org.thingsboard.server.common.data.notification.NotificationRequest; import org.thingsboard.server.common.data.notification.NotificationRequestConfig; +import org.thingsboard.server.common.data.notification.NotificationRequestStatus; import org.thingsboard.server.common.data.notification.rule.NonConfirmedNotificationEscalation; import org.thingsboard.server.common.data.notification.rule.NotificationRule; import org.thingsboard.server.common.data.notification.rule.NotificationRuleConfig; +import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; +import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg; import org.thingsboard.server.dao.notification.NotificationRequestService; import org.thingsboard.server.dao.notification.NotificationRuleService; import org.thingsboard.server.queue.util.TbCoreComponent; -import org.thingsboard.server.service.executors.DbCallbackExecutorService; +import org.thingsboard.server.service.executors.NotificationExecutorService; import java.util.List; import java.util.Map; @@ -53,7 +59,7 @@ public class DefaultNotificationRuleProcessingService implements NotificationRul private final NotificationRequestService notificationRequestService; @Autowired @Lazy private NotificationManager notificationManager; - private final DbCallbackExecutorService dbCallbackExecutorService; + private final NotificationExecutorService notificationExecutor; @Override public ListenableFuture onAlarmCreatedOrUpdated(TenantId tenantId, Alarm alarm) { @@ -67,23 +73,12 @@ public class DefaultNotificationRuleProcessingService implements NotificationRul private ListenableFuture processAlarmUpdate(TenantId tenantId, Alarm alarm, boolean deleted) { if (alarm.getNotificationRuleId() == null) return Futures.immediateFuture(null); - return dbCallbackExecutorService.submit(() -> { + return notificationExecutor.submit(() -> { onAlarmUpdate(tenantId, alarm.getNotificationRuleId(), alarm, deleted); return null; }); } - @Override - public ListenableFuture onNotificationRuleDeleted(TenantId tenantId, NotificationRuleId ruleId) { - return dbCallbackExecutorService.submit(() -> { - // FIXME [viacheslav] - // need to remove fk constraint in notificationRequest to rule - // todo: do we need to remove all notifications when notification request is deleted? - return null; - }); - } - - // todo: think about: what if notification rule was updated? private void onAlarmUpdate(TenantId tenantId, NotificationRuleId notificationRuleId, Alarm alarm, boolean deleted) { log.debug("Processing alarm update ({}) with notification rule {}", alarm.getId(), notificationRuleId); List notificationRequests = notificationRequestService.findNotificationRequestsByRuleIdAndOriginatorEntityId(tenantId, notificationRuleId, alarm.getId()); @@ -91,11 +86,14 @@ public class DefaultNotificationRuleProcessingService implements NotificationRul if (notificationRule == null) return; if (alarmAcknowledged(alarm) || deleted) { + if (notificationRequests.isEmpty()) { + return; + } for (NotificationRequest notificationRequest : notificationRequests) { - notificationManager.deleteNotificationRequest(tenantId, notificationRequest.getId()); - // todo: or should we mark already sent notifications as read and delete only scheduled? + if (notificationRequest.getStatus() == NotificationRequestStatus.SCHEDULED) { + notificationManager.deleteNotificationRequest(tenantId, notificationRequest.getId()); + } } - return; } if (notificationRequests.isEmpty()) { @@ -139,20 +137,22 @@ public class DefaultNotificationRuleProcessingService implements NotificationRul ); NotificationRequest notificationRequest = NotificationRequest.builder() .tenantId(tenantId) - .targetId(targetId) + .targets(List.of(targetId)) .templateId(notificationRule.getTemplateId()) .deliveryMethods(notificationRule.getDeliveryMethods()) .additionalConfig(config) .info(notificationInfo) + .ruleId(notificationRule.getId()) .originatorType(NotificationOriginatorType.ALARM) .originatorEntityId(alarm.getId()) - .ruleId(notificationRule.getId()) + .originatorEntity(alarm) .templateContext(templateContext) .build(); notificationManager.processNotificationRequest(tenantId, notificationRequest); } private NotificationInfo constructNotificationInfo(Alarm alarm) { + // TODO: add info about assignee return AlarmOriginatedNotificationInfo.builder() .alarmId(alarm.getId()) .alarmType(alarm.getType()) @@ -162,4 +162,19 @@ public class DefaultNotificationRuleProcessingService implements NotificationRul .build(); } + @EventListener(ComponentLifecycleMsg.class) + public void onNotificationRuleDeleted(ComponentLifecycleMsg componentLifecycleMsg) { + if (componentLifecycleMsg.getEvent() != ComponentLifecycleEvent.DELETED || + componentLifecycleMsg.getEntityId().getEntityType() != EntityType.NOTIFICATION_RULE) { + return; + } + + TenantId tenantId = componentLifecycleMsg.getTenantId(); + NotificationRuleId notificationRuleId = (NotificationRuleId) componentLifecycleMsg.getEntityId(); + List scheduledForRule = notificationRequestService.findNotificationRequestsIdsByStatusAndRuleId(tenantId, NotificationRequestStatus.SCHEDULED, notificationRuleId); + for (NotificationRequestId notificationRequestId : scheduledForRule) { + notificationManager.deleteNotificationRequest(tenantId, notificationRequestId); + } + } + } diff --git a/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationSchedulerService.java b/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationSchedulerService.java index a2757dc3dd..67ffa8e167 100644 --- a/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationSchedulerService.java +++ b/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationSchedulerService.java @@ -16,24 +16,31 @@ package org.thingsboard.server.service.notification; import com.google.common.util.concurrent.ListenableFuture; -import com.google.common.util.concurrent.ListenableScheduledFuture; +import lombok.Data; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; +import org.springframework.context.event.EventListener; import org.springframework.stereotype.Service; import org.thingsboard.rule.engine.api.NotificationManager; +import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.NotificationRequestId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.notification.NotificationRequest; import org.thingsboard.server.common.data.notification.NotificationRequestConfig; import org.thingsboard.server.common.data.page.PageDataIterable; +import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; +import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.dao.notification.NotificationRequestService; +import org.thingsboard.server.queue.scheduler.SchedulerComponent; import org.thingsboard.server.queue.util.TbCoreComponent; +import org.thingsboard.server.service.executors.NotificationExecutorService; import org.thingsboard.server.service.partition.AbstractPartitionBasedService; import javax.annotation.PostConstruct; import java.util.Collections; +import java.util.HashSet; import java.util.List; import java.util.Map; import java.util.Optional; @@ -51,8 +58,10 @@ public class DefaultNotificationSchedulerService extends AbstractPartitionBasedS private final NotificationManager notificationManager; private final NotificationRequestService notificationRequestService; + private final SchedulerComponent scheduler; + private final NotificationExecutorService notificationExecutor; - private final Map> scheduledNotificationRequests = new ConcurrentHashMap<>(); + private final Map scheduledNotificationRequests = new ConcurrentHashMap<>(); @PostConstruct public void init() { @@ -93,30 +102,48 @@ public class DefaultNotificationSchedulerService extends AbstractPartitionBasedS delayInMs = 0; } - ListenableScheduledFuture scheduledTask = scheduledExecutor.schedule(() -> { + ScheduledFuture scheduledTask = scheduler.schedule(() -> { NotificationRequest notificationRequest = notificationRequestService.findNotificationRequestById(tenantId, request.getId()); if (notificationRequest == null) return; - notificationManager.processNotificationRequest(tenantId, notificationRequest); + notificationExecutor.executeAsync(() -> { + notificationManager.processNotificationRequest(tenantId, notificationRequest); + }); scheduledNotificationRequests.remove(notificationRequest.getId()); }, delayInMs, TimeUnit.MILLISECONDS); - scheduledNotificationRequests.put(request.getId(), scheduledTask); + scheduledNotificationRequests.put(request.getId(), new ScheduledRequestMetadata(tenantId, scheduledTask)); } - @Override - public void onNotificationRequestDeleted(TenantId tenantId, NotificationRequestId notificationRequestId) { - removeAndCancel(notificationRequestId); + @EventListener(ComponentLifecycleMsg.class) + public void handleComponentLifecycleEvent(ComponentLifecycleMsg event) { + if (event.getEvent() == ComponentLifecycleEvent.DELETED) { + EntityId entityId = event.getEntityId(); + switch (entityId.getEntityType()) { + case NOTIFICATION_REQUEST: + cancelAndRemove((NotificationRequestId) entityId); + break; + case TENANT: + Set toCancel = new HashSet<>(); + scheduledNotificationRequests.forEach((notificationRequestId, scheduledRequestMetadata) -> { + if (scheduledRequestMetadata.getTenantId().equals(entityId)) { + toCancel.add(notificationRequestId); + } + }); + toCancel.forEach(this::cancelAndRemove); + break; + } + } } @Override protected void cleanupEntityOnPartitionRemoval(NotificationRequestId notificationRequestId) { - removeAndCancel(notificationRequestId); + cancelAndRemove(notificationRequestId); } - private void removeAndCancel(NotificationRequestId notificationRequestId) { - ScheduledFuture scheduledTask = scheduledNotificationRequests.remove(notificationRequestId); - if (scheduledTask != null) { - scheduledTask.cancel(false); + private void cancelAndRemove(NotificationRequestId notificationRequestId) { + ScheduledRequestMetadata md = scheduledNotificationRequests.remove(notificationRequestId); + if (md != null) { + md.getFuture().cancel(false); } } @@ -130,4 +157,10 @@ public class DefaultNotificationSchedulerService extends AbstractPartitionBasedS return "notifications-scheduler"; } + @Data + private static class ScheduledRequestMetadata { + private final TenantId tenantId; + private final ScheduledFuture future; + } + } diff --git a/application/src/main/java/org/thingsboard/server/service/notification/NotificationProcessingContext.java b/application/src/main/java/org/thingsboard/server/service/notification/NotificationProcessingContext.java index 9abcba6b43..46cc211403 100644 --- a/application/src/main/java/org/thingsboard/server/service/notification/NotificationProcessingContext.java +++ b/application/src/main/java/org/thingsboard/server/service/notification/NotificationProcessingContext.java @@ -18,7 +18,10 @@ package org.thingsboard.server.service.notification; import com.google.common.base.Strings; import lombok.Builder; import lombok.Getter; +import org.apache.commons.lang3.StringUtils; +import org.thingsboard.server.common.data.HasCustomerId; import org.thingsboard.server.common.data.User; +import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.notification.NotificationDeliveryMethod; import org.thingsboard.server.common.data.notification.NotificationRequest; @@ -32,14 +35,13 @@ import org.thingsboard.server.dao.notification.NotificationTemplateService; import java.util.HashMap; import java.util.Map; -import java.util.Optional; -import java.util.stream.Collectors; @SuppressWarnings("unchecked") public class NotificationProcessingContext { @Getter private final TenantId tenantId; + private final HasCustomerId originatorEntity; private final NotificationSettings settings; @Getter private final NotificationRequest request; @@ -49,26 +51,28 @@ public class NotificationProcessingContext { private NotificationTemplate notificationTemplate; private Map templates; @Getter - private NotificationRequestStats stats; + private final NotificationRequestStats stats; @Builder - public NotificationProcessingContext(TenantId tenantId, NotificationSettings settings, NotificationRequest request, Map additionalTemplateContext) { + public NotificationProcessingContext(TenantId tenantId, NotificationRequest request, NotificationSettings settings) { this.tenantId = tenantId; - this.settings = settings; + this.originatorEntity = request.getOriginatorEntity(); this.request = request; - this.additionalTemplateContext = additionalTemplateContext; + this.settings = settings; + this.additionalTemplateContext = request.getTemplateContext(); + this.stats = new NotificationRequestStats(); } public void init(NotificationTemplateService templateService) { notificationTemplate = templateService.findNotificationTemplateById(tenantId, request.getTemplateId()); NotificationTemplateConfig config = notificationTemplate.getConfiguration(); - templates = request.getDeliveryMethods().stream() - .collect(Collectors.toMap(k -> k, deliveryMethod -> { - return Optional.ofNullable(config.getTemplates()) - .map(templates -> templates.get(deliveryMethod)) - .orElse(config.getDefaultTemplate()); - })); - stats = new NotificationRequestStats(); + for (NotificationDeliveryMethod deliveryMethod : request.getDeliveryMethods()) { + DeliveryMethodNotificationTemplate template = config.getTemplates().get(deliveryMethod); + if (StringUtils.isEmpty(template.getBody())) { + template.setBody(config.getDefaultTextTemplate()); + } + } + templates = config.getTemplates(); } public T getTemplate(NotificationDeliveryMethod deliveryMethod) { @@ -90,4 +94,8 @@ public class NotificationProcessingContext { return templateContext; } + public CustomerId getOriginatorCustomerId() { + return originatorEntity != null ? originatorEntity.getCustomerId() : null; + } + } diff --git a/application/src/main/java/org/thingsboard/server/service/notification/NotificationRuleProcessingService.java b/application/src/main/java/org/thingsboard/server/service/notification/NotificationRuleProcessingService.java index 2283e49cdd..8eb550e12a 100644 --- a/application/src/main/java/org/thingsboard/server/service/notification/NotificationRuleProcessingService.java +++ b/application/src/main/java/org/thingsboard/server/service/notification/NotificationRuleProcessingService.java @@ -17,7 +17,6 @@ package org.thingsboard.server.service.notification; import com.google.common.util.concurrent.ListenableFuture; import org.thingsboard.server.common.data.alarm.Alarm; -import org.thingsboard.server.common.data.id.NotificationRuleId; import org.thingsboard.server.common.data.id.TenantId; public interface NotificationRuleProcessingService { @@ -26,6 +25,4 @@ public interface NotificationRuleProcessingService { ListenableFuture onAlarmDeleted(TenantId tenantId, Alarm alarm); - ListenableFuture onNotificationRuleDeleted(TenantId tenantId, NotificationRuleId ruleId); - } diff --git a/application/src/main/java/org/thingsboard/server/service/notification/NotificationSchedulerService.java b/application/src/main/java/org/thingsboard/server/service/notification/NotificationSchedulerService.java index ce018f171c..01d3aea1d3 100644 --- a/application/src/main/java/org/thingsboard/server/service/notification/NotificationSchedulerService.java +++ b/application/src/main/java/org/thingsboard/server/service/notification/NotificationSchedulerService.java @@ -22,6 +22,4 @@ public interface NotificationSchedulerService { void scheduleNotificationRequest(TenantId tenantId, NotificationRequestId notificationRequestId, long requestTs); - void onNotificationRequestDeleted(TenantId tenantId, NotificationRequestId notificationRequestId); - } diff --git a/application/src/main/java/org/thingsboard/server/service/notification/channels/SlackNotificationChannel.java b/application/src/main/java/org/thingsboard/server/service/notification/channels/SlackNotificationChannel.java index a12efbffa0..525feb1a40 100644 --- a/application/src/main/java/org/thingsboard/server/service/notification/channels/SlackNotificationChannel.java +++ b/application/src/main/java/org/thingsboard/server/service/notification/channels/SlackNotificationChannel.java @@ -20,7 +20,7 @@ import com.google.common.util.concurrent.ListenableFuture; import lombok.RequiredArgsConstructor; import org.apache.commons.lang3.StringUtils; import org.springframework.stereotype.Component; -import org.thingsboard.rule.engine.api.slack.SlackConversation; +import org.thingsboard.server.common.data.notification.template.SlackConversation; import org.thingsboard.rule.engine.api.slack.SlackService; import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.notification.AlreadySentException; diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java index d71613cab4..09441f2563 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java @@ -21,6 +21,7 @@ import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.StringUtils; import org.springframework.beans.factory.annotation.Value; import org.springframework.boot.context.event.ApplicationReadyEvent; +import org.springframework.context.ApplicationEventPublisher; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Service; import org.thingsboard.common.util.JacksonUtil; @@ -38,8 +39,6 @@ import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.common.msg.rpc.FromDeviceRpcResponse; import org.thingsboard.server.common.stats.StatsFactory; -import org.thingsboard.server.service.security.auth.jwt.settings.JwtSettingsService; -import org.thingsboard.server.queue.util.DataDecodingEncodingService; import org.thingsboard.server.dao.tenant.TbTenantProfileCache; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.DeviceStateServiceMsgProto; @@ -65,6 +64,7 @@ import org.thingsboard.server.queue.discovery.PartitionService; import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent; import org.thingsboard.server.queue.provider.TbCoreQueueFactory; import org.thingsboard.server.queue.util.AfterStartUp; +import org.thingsboard.server.queue.util.DataDecodingEncodingService; import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.apiusage.TbApiUsageStateService; import org.thingsboard.server.service.edge.EdgeNotificationService; @@ -76,6 +76,7 @@ import org.thingsboard.server.service.queue.processing.AbstractConsumerService; import org.thingsboard.server.service.queue.processing.IdMsgPair; import org.thingsboard.server.service.rpc.TbCoreDeviceRpcService; import org.thingsboard.server.service.rpc.ToDeviceRpcRequestActorMsg; +import org.thingsboard.server.service.security.auth.jwt.settings.JwtSettingsService; import org.thingsboard.server.service.state.DeviceStateService; import org.thingsboard.server.service.subscription.SubscriptionManagerService; import org.thingsboard.server.service.subscription.TbLocalSubscriptionService; @@ -153,9 +154,10 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService jwtSettingsService, NotificationSchedulerService notificationSchedulerService) { - super(actorContext, encodingService, tenantProfileCache, deviceProfileCache, assetProfileCache, apiUsageStateService, partitionService, tbCoreQueueFactory.createToCoreNotificationsMsgConsumer(), jwtSettingsService); + super(actorContext, encodingService, tenantProfileCache, deviceProfileCache, assetProfileCache, apiUsageStateService, partitionService, eventPublisher, tbCoreQueueFactory.createToCoreNotificationsMsgConsumer(), jwtSettingsService); this.mainConsumer = tbCoreQueueFactory.createToCoreMsgConsumer(); this.usageStatsConsumer = tbCoreQueueFactory.createToUsageStatsServiceMsgConsumer(); this.firmwareStatesConsumer = tbCoreQueueFactory.createToOtaPackageStateServiceMsgConsumer(); @@ -584,13 +586,8 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService> nfConsumer; protected final Optional jwtSettingsService; @@ -83,7 +85,8 @@ public abstract class AbstractConsumerService> nfConsumer, Optional jwtSettingsService) { + PartitionService partitionService, ApplicationEventPublisher eventPublisher, + TbQueueConsumer> nfConsumer, Optional jwtSettingsService) { this.actorContext = actorContext; this.encodingService = encodingService; this.tenantProfileCache = tenantProfileCache; @@ -91,6 +94,7 @@ public abstract class AbstractConsumerService 0) { + long gap = TimeUnit.MINUTES.toMillis(10); + long requestExpTime = lastRemovedNotificationTs - TimeUnit.SECONDS.toMillis(NotificationRequestConfig.MAX_SENDING_DELAY) - gap; + // TODO: double-check this + notificationRequestDao.removeAllByCreatedTimeBefore(requestExpTime); } } diff --git a/application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java b/application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java index 9ea269918e..6d4a5bcb1e 100644 --- a/application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java @@ -296,21 +296,24 @@ public abstract class AbstractWebTest extends AbstractInMemoryStorageTest { if (savedDifferentTenant != null) { login(DIFFERENT_TENANT_ADMIN_EMAIL, DIFFERENT_TENANT_ADMIN_PASSWORD); } else { - loginSysAdmin(); - - Tenant tenant = new Tenant(); - tenant.setTitle(TEST_DIFFERENT_TENANT_NAME); - savedDifferentTenant = doPost("/api/tenant", tenant, Tenant.class); - differentTenantId = savedDifferentTenant.getId(); - Assert.assertNotNull(savedDifferentTenant); - User differentTenantAdmin = new User(); - differentTenantAdmin.setAuthority(Authority.TENANT_ADMIN); - differentTenantAdmin.setTenantId(savedDifferentTenant.getId()); - differentTenantAdmin.setEmail(DIFFERENT_TENANT_ADMIN_EMAIL); - savedDifferentTenantUser = createUserAndLogin(differentTenantAdmin, DIFFERENT_TENANT_ADMIN_PASSWORD); + createDifferentTenant(); } } + protected void createDifferentTenant() throws Exception { + loginSysAdmin(); + Tenant tenant = new Tenant(); + tenant.setTitle(TEST_DIFFERENT_TENANT_NAME); + savedDifferentTenant = doPost("/api/tenant", tenant, Tenant.class); + differentTenantId = savedDifferentTenant.getId(); + Assert.assertNotNull(savedDifferentTenant); + User differentTenantAdmin = new User(); + differentTenantAdmin.setAuthority(Authority.TENANT_ADMIN); + differentTenantAdmin.setTenantId(savedDifferentTenant.getId()); + differentTenantAdmin.setEmail(DIFFERENT_TENANT_ADMIN_EMAIL); + savedDifferentTenantUser = createUserAndLogin(differentTenantAdmin, DIFFERENT_TENANT_ADMIN_PASSWORD); + } + protected void loginDifferentCustomer() throws Exception { if (savedDifferentCustomer != null) { login(savedDifferentCustomer.getEmail(), CUSTOMER_USER_PASSWORD); diff --git a/application/src/test/java/org/thingsboard/server/service/notification/NotificationApiTest.java b/application/src/test/java/org/thingsboard/server/service/notification/NotificationApiTest.java index e7acd85fab..1793ccb4b4 100644 --- a/application/src/test/java/org/thingsboard/server/service/notification/NotificationApiTest.java +++ b/application/src/test/java/org/thingsboard/server/service/notification/NotificationApiTest.java @@ -34,8 +34,6 @@ import org.thingsboard.server.common.data.notification.NotificationRequest; import org.thingsboard.server.common.data.notification.NotificationRequestConfig; import org.thingsboard.server.common.data.notification.NotificationRequestStats; import org.thingsboard.server.common.data.notification.NotificationRequestStatus; -import org.thingsboard.server.common.data.notification.settings.NotificationDeliveryMethodConfig; -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.SingleUserNotificationTargetConfig; import org.thingsboard.server.common.data.notification.targets.UserListNotificationTargetConfig; @@ -366,6 +364,9 @@ public class NotificationApiTest extends AbstractControllerTest { } private NotificationRequest submitNotificationRequest(NotificationTargetId targetId, String text, int delayInSec, NotificationDeliveryMethod... deliveryMethods) { + if (deliveryMethods.length == 0) { + deliveryMethods = new NotificationDeliveryMethod[]{NotificationDeliveryMethod.WEBSOCKET}; + } NotificationTemplate notificationTemplate = createNotificationTemplate(text, deliveryMethods); NotificationRequestConfig config = new NotificationRequestConfig(); config.setSendingDelayInSec(delayInSec); @@ -373,10 +374,10 @@ public class NotificationApiTest extends AbstractControllerTest { notificationInfo.setDescription("The text: " + text); NotificationRequest notificationRequest = NotificationRequest.builder() .tenantId(tenantId) - .targetId(targetId) + .targets(List.of(targetId)) .templateId(notificationTemplate.getId()) .info(notificationInfo) - .deliveryMethods(deliveryMethods.length > 0 ? List.of(deliveryMethods) : List.of(NotificationDeliveryMethod.WEBSOCKET)) + .deliveryMethods(List.of(deliveryMethods)) .additionalConfig(config) .build(); return doPost("/api/notification/request", notificationRequest, NotificationRequest.class); @@ -388,18 +389,18 @@ public class NotificationApiTest extends AbstractControllerTest { notificationTemplate.setName("Notification template for testing"); notificationTemplate.setNotificationType("Just a test"); NotificationTemplateConfig config = new NotificationTemplateConfig(); - DeliveryMethodNotificationTemplate defaultTemplate = new DeliveryMethodNotificationTemplate(); - defaultTemplate.setBody(text); - config.setDefaultTemplate(defaultTemplate); + config.setDefaultTextTemplate(text); + config.setTemplates(new HashMap<>()); for (NotificationDeliveryMethod deliveryMethod : deliveryMethods) { if (deliveryMethod == NotificationDeliveryMethod.EMAIL) { EmailDeliveryMethodNotificationTemplate emailNotificationTemplate = new EmailDeliveryMethodNotificationTemplate(); emailNotificationTemplate.setSubject("Hello from test"); - emailNotificationTemplate.setBody(text); emailNotificationTemplate.setMethod(deliveryMethod); - config.setTemplates(Map.of( - deliveryMethod, emailNotificationTemplate - )); + config.getTemplates().put(deliveryMethod, emailNotificationTemplate); + } else { + DeliveryMethodNotificationTemplate defaultTemplate = new DeliveryMethodNotificationTemplate(); + defaultTemplate.setMethod(deliveryMethod); + config.getTemplates().put(deliveryMethod, defaultTemplate); } } notificationTemplate.setConfiguration(config); diff --git a/application/src/test/java/org/thingsboard/server/service/notification/NotificationTargetApiTest.java b/application/src/test/java/org/thingsboard/server/service/notification/NotificationTargetApiTest.java new file mode 100644 index 0000000000..d9bcd0f6e2 --- /dev/null +++ b/application/src/test/java/org/thingsboard/server/service/notification/NotificationTargetApiTest.java @@ -0,0 +1,168 @@ +/** + * Copyright © 2016-2022 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; + +import com.fasterxml.jackson.core.type.TypeReference; +import org.junit.Before; +import org.junit.Test; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.test.web.servlet.ResultActions; +import org.springframework.test.web.servlet.ResultMatcher; +import org.thingsboard.server.common.data.User; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.notification.targets.AllUsersNotificationTargetConfig; +import org.thingsboard.server.common.data.notification.targets.CustomerUsersNotificationTargetConfig; +import org.thingsboard.server.common.data.notification.targets.NotificationTarget; +import org.thingsboard.server.common.data.notification.targets.SingleUserNotificationTargetConfig; +import org.thingsboard.server.common.data.notification.targets.UserListNotificationTargetConfig; +import org.thingsboard.server.common.data.page.PageData; +import org.thingsboard.server.controller.AbstractControllerTest; +import org.thingsboard.server.dao.notification.NotificationTargetDao; +import org.thingsboard.server.dao.service.DaoSqlTest; + +import java.util.Collections; +import java.util.List; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; + +@DaoSqlTest +public class NotificationTargetApiTest extends AbstractControllerTest { + + @Autowired + private NotificationTargetDao notificationTargetDao; + + @Before + public void beforeEach() throws Exception { + loginTenantAdmin(); + } + + @Test + public void givenInvalidNotificationTarget_whenSaving_returnValidationError() throws Exception { + NotificationTarget notificationTarget = new NotificationTarget(); + notificationTarget.setTenantId(null); + notificationTarget.setName(null); + notificationTarget.setConfiguration(null); + + String validationError = saveAndGetError(notificationTarget, status().isBadRequest()); + assertThat(validationError) + .contains("name must not be") + .contains("configuration must not be"); + + SingleUserNotificationTargetConfig singleUserConfig = new SingleUserNotificationTargetConfig(); + singleUserConfig.setUserId(null); + notificationTarget.setConfiguration(singleUserConfig); + + validationError = saveAndGetError(notificationTarget, status().isBadRequest()); + assertThat(validationError) + .contains("userId must not be"); + + UserListNotificationTargetConfig userListConfig = new UserListNotificationTargetConfig(); + userListConfig.setUsersIds(Collections.emptyList()); + notificationTarget.setConfiguration(userListConfig); + + validationError = saveAndGetError(notificationTarget, status().isBadRequest()); + assertThat(validationError) + .contains("usersIds must not be"); + } + + @Test + public void givenNotificationTargetWithUsersFromDifferentTenant_whenSaving_returnAccessDeniedError() throws Exception { + loginDifferentTenant(); + NotificationTarget notificationTarget = new NotificationTarget(); + notificationTarget.setTenantId(differentTenantId); + notificationTarget.setName("Target 1"); + UserListNotificationTargetConfig userListConfig = new UserListNotificationTargetConfig(); + userListConfig.setUsersIds(List.of(customerUserId.getId(), tenantAdminUserId.getId())); + notificationTarget.setConfiguration(userListConfig); + + saveAndGetError(notificationTarget, status().isForbidden()); + + SingleUserNotificationTargetConfig singleUserConfig = new SingleUserNotificationTargetConfig(); + singleUserConfig.setUserId(customerUserId.getId()); + notificationTarget.setConfiguration(singleUserConfig); + + saveAndGetError(notificationTarget, status().isForbidden()); + + loginSysAdmin(); + notificationTarget.setTenantId(TenantId.SYS_TENANT_ID); + notificationTarget.setConfiguration(userListConfig); + save(notificationTarget, status().isOk()); + } + + @Test + public void givenNotificationTargetConfig_testGetRecipients() throws Exception { + NotificationTarget notificationTarget = new NotificationTarget(); + notificationTarget.setTenantId(tenantId); + notificationTarget.setName("Test target"); + CustomerUsersNotificationTargetConfig customerUsersConfig = new CustomerUsersNotificationTargetConfig(); + customerUsersConfig.setCustomerId(customerId.getId()); + notificationTarget.setConfiguration(customerUsersConfig); + + List recipients = getRecipients(notificationTarget); + assertThat(recipients).size().isNotZero(); + assertThat(recipients).allSatisfy(recipient -> { + assertThat(recipient.getCustomerId()).isEqualTo(customerId); + }); + + AllUsersNotificationTargetConfig allUsersConfig = new AllUsersNotificationTargetConfig(); + notificationTarget.setConfiguration(allUsersConfig); + recipients = getRecipients(notificationTarget); + assertThat(recipients).size().isGreaterThanOrEqualTo(2); + assertThat(recipients).allSatisfy(recipient -> { + assertThat(recipient.getTenantId()).isEqualTo(tenantId); + }); + + createDifferentTenant(); + loginSysAdmin(); + recipients = getRecipients(notificationTarget); + assertThat(recipients).size().isGreaterThanOrEqualTo(3); + assertThat(recipients).anySatisfy(recipient -> { + assertThat(recipient.getTenantId()).isEqualTo(tenantId); + }); + assertThat(recipients).anySatisfy(recipient -> { + assertThat(recipient.getTenantId()).isEqualTo(differentTenantId); + }); + } + + @Test + public void whenDeletingTenant_thenDeleteNotificationTarget() throws Exception { + createDifferentTenant(); + NotificationTarget notificationTarget = new NotificationTarget(); + notificationTarget.setName("Test 1"); + notificationTarget.setTenantId(differentTenantId); + notificationTarget.setConfiguration(new AllUsersNotificationTargetConfig()); + save(notificationTarget, status().isOk()); + assertThat(notificationTargetDao.find(TenantId.SYS_TENANT_ID)).isNotEmpty(); + + deleteDifferentTenant(); + assertThat(notificationTargetDao.find(TenantId.SYS_TENANT_ID)).isEmpty(); + } + + private String saveAndGetError(NotificationTarget notificationTarget, ResultMatcher statusMatcher) throws Exception { + return getErrorMessage(save(notificationTarget, statusMatcher)); + } + + private ResultActions save(NotificationTarget notificationTarget, ResultMatcher statusMatcher) throws Exception { + return doPost("/api/notification/target", notificationTarget) + .andExpect(statusMatcher); + } + + private List getRecipients(NotificationTarget notificationTarget) throws Exception { + return doPostWithTypedResponse("/api/notification/target/recipients?page=0&pageSize=100", notificationTarget, new TypeReference>() {}).getData(); + } + +} diff --git a/application/src/test/java/org/thingsboard/server/service/notification/NotificationTemplateApiTest.java b/application/src/test/java/org/thingsboard/server/service/notification/NotificationTemplateApiTest.java new file mode 100644 index 0000000000..58135472ab --- /dev/null +++ b/application/src/test/java/org/thingsboard/server/service/notification/NotificationTemplateApiTest.java @@ -0,0 +1,90 @@ +/** + * Copyright © 2016-2022 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; + +import org.junit.Before; +import org.junit.Test; +import org.springframework.test.web.servlet.ResultActions; +import org.springframework.test.web.servlet.ResultMatcher; +import org.thingsboard.server.common.data.notification.NotificationDeliveryMethod; +import org.thingsboard.server.common.data.notification.template.EmailDeliveryMethodNotificationTemplate; +import org.thingsboard.server.common.data.notification.template.NotificationTemplate; +import org.thingsboard.server.common.data.notification.template.NotificationTemplateConfig; +import org.thingsboard.server.controller.AbstractControllerTest; +import org.thingsboard.server.dao.service.DaoSqlTest; + +import java.util.Map; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; + +@DaoSqlTest +public class NotificationTemplateApiTest extends AbstractControllerTest { + + @Before + public void beforeEach() throws Exception { + loginTenantAdmin(); + } + + @Test + public void givenInvalidNotificationTemplate_whenSaving_returnValidationError() throws Exception { + NotificationTemplate notificationTemplate = new NotificationTemplate(); + notificationTemplate.setTenantId(tenantId); + notificationTemplate.setName(null); + notificationTemplate.setNotificationType(null); + notificationTemplate.setConfiguration(null); + + String validationError = saveAndGetError(notificationTemplate, status().isBadRequest()); + assertThat(validationError) + .contains("name must not be") + .contains("notificationType must not be") + .contains("configuration must not be"); + + NotificationTemplateConfig config = new NotificationTemplateConfig(); + notificationTemplate.setConfiguration(config); + config.setDefaultTextTemplate("Default text"); + EmailDeliveryMethodNotificationTemplate emailTemplate = new EmailDeliveryMethodNotificationTemplate(); + emailTemplate.setMethod(NotificationDeliveryMethod.EMAIL); + emailTemplate.setBody(null); + emailTemplate.setSubject(null); + config.setTemplates(Map.of( + NotificationDeliveryMethod.EMAIL, emailTemplate + )); + notificationTemplate.setName("