Browse Source

Multiple notification targets for request; improvements and refactoring

pull/7511/head
ViacheslavKlimov 4 years ago
parent
commit
409976c6a3
  1. 14
      application/src/main/data/upgrade/3.4.3/schema_update.sql
  2. 5
      application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java
  3. 5
      application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java
  4. 12
      application/src/main/java/org/thingsboard/server/controller/BaseController.java
  5. 3
      application/src/main/java/org/thingsboard/server/controller/NotificationController.java
  6. 7
      application/src/main/java/org/thingsboard/server/controller/NotificationRuleController.java
  7. 20
      application/src/main/java/org/thingsboard/server/controller/NotificationTargetController.java
  8. 3
      application/src/main/java/org/thingsboard/server/controller/NotificationTemplateController.java
  9. 29
      application/src/main/java/org/thingsboard/server/service/executors/NotificationExecutorService.java
  10. 5
      application/src/main/java/org/thingsboard/server/service/install/DefaultSystemDataLoaderService.java
  11. 3
      application/src/main/java/org/thingsboard/server/service/install/SystemDataLoaderService.java
  12. 80
      application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationManager.java
  13. 53
      application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationRuleProcessingService.java
  14. 59
      application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationSchedulerService.java
  15. 34
      application/src/main/java/org/thingsboard/server/service/notification/NotificationProcessingContext.java
  16. 3
      application/src/main/java/org/thingsboard/server/service/notification/NotificationRuleProcessingService.java
  17. 2
      application/src/main/java/org/thingsboard/server/service/notification/NotificationSchedulerService.java
  18. 2
      application/src/main/java/org/thingsboard/server/service/notification/channels/SlackNotificationChannel.java
  19. 15
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java
  20. 6
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java
  21. 7
      application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java
  22. 2
      application/src/main/java/org/thingsboard/server/service/slack/DefaultSlackService.java
  23. 20
      application/src/main/java/org/thingsboard/server/service/ttl/NotificationsCleanUpService.java
  24. 27
      application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java
  25. 23
      application/src/test/java/org/thingsboard/server/service/notification/NotificationApiTest.java
  26. 168
      application/src/test/java/org/thingsboard/server/service/notification/NotificationTargetApiTest.java
  27. 90
      application/src/test/java/org/thingsboard/server/service/notification/NotificationTemplateApiTest.java
  28. 86
      application/src/test/java/org/thingsboard/server/service/notification/NotificationsClient.java
  29. 1
      common/cluster-api/src/main/proto/queue.proto
  30. 3
      common/dao-api/src/main/java/org/thingsboard/server/dao/notification/NotificationRequestService.java
  31. 4
      common/dao-api/src/main/java/org/thingsboard/server/dao/notification/NotificationRuleService.java
  32. 7
      common/dao-api/src/main/java/org/thingsboard/server/dao/notification/NotificationTargetService.java
  33. 14
      common/data/src/main/java/org/thingsboard/server/common/data/notification/NotificationRequest.java
  34. 4
      common/data/src/main/java/org/thingsboard/server/common/data/notification/NotificationRequestConfig.java
  35. 15
      common/data/src/main/java/org/thingsboard/server/common/data/notification/NotificationRequestStats.java
  36. 5
      common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/NonConfirmedNotificationEscalation.java
  37. 16
      common/data/src/main/java/org/thingsboard/server/common/data/notification/targets/CustomerUsersNotificationTargetConfig.java
  38. 3
      common/data/src/main/java/org/thingsboard/server/common/data/notification/targets/NotificationTarget.java
  39. 3
      common/data/src/main/java/org/thingsboard/server/common/data/notification/template/DeliveryMethodNotificationTemplate.java
  40. 3
      common/data/src/main/java/org/thingsboard/server/common/data/notification/template/EmailDeliveryMethodNotificationTemplate.java
  41. 1
      common/data/src/main/java/org/thingsboard/server/common/data/notification/template/NotificationTemplate.java
  42. 19
      common/data/src/main/java/org/thingsboard/server/common/data/notification/template/NotificationTemplateConfig.java
  43. 2
      common/data/src/main/java/org/thingsboard/server/common/data/notification/template/SlackConversation.java
  44. 2
      common/data/src/main/java/org/thingsboard/server/common/data/notification/template/SlackDeliveryMethodNotificationTemplate.java
  45. 44
      dao/src/main/java/org/thingsboard/server/dao/model/BaseSqlEntity.java
  46. 2
      dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java
  47. 2
      dao/src/main/java/org/thingsboard/server/dao/model/sql/AbstractAlarmEntity.java
  48. 5
      dao/src/main/java/org/thingsboard/server/dao/model/sql/NotificationEntity.java
  49. 27
      dao/src/main/java/org/thingsboard/server/dao/model/sql/NotificationRequestEntity.java
  50. 17
      dao/src/main/java/org/thingsboard/server/dao/model/sql/NotificationRuleEntity.java
  51. 4
      dao/src/main/java/org/thingsboard/server/dao/model/sql/NotificationTargetEntity.java
  52. 4
      dao/src/main/java/org/thingsboard/server/dao/model/sql/NotificationTemplateEntity.java
  53. 9
      dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationRequestService.java
  54. 8
      dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationRuleService.java
  55. 38
      dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationTargetService.java
  56. 19
      dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationTemplateService.java
  57. 10
      dao/src/main/java/org/thingsboard/server/dao/notification/NotificationRequestDao.java
  58. 3
      dao/src/main/java/org/thingsboard/server/dao/notification/NotificationRuleDao.java
  59. 27
      dao/src/main/java/org/thingsboard/server/dao/service/ConstraintValidator.java
  60. 1
      dao/src/main/java/org/thingsboard/server/dao/service/NoXssValidator.java
  61. 24
      dao/src/main/java/org/thingsboard/server/dao/sql/notification/JpaNotificationRequestDao.java
  62. 6
      dao/src/main/java/org/thingsboard/server/dao/sql/notification/JpaNotificationRuleDao.java
  63. 10
      dao/src/main/java/org/thingsboard/server/dao/sql/notification/NotificationRequestRepository.java
  64. 2
      dao/src/main/java/org/thingsboard/server/dao/sql/notification/NotificationRuleRepository.java
  65. 5
      dao/src/main/java/org/thingsboard/server/dao/sqlts/insert/sql/SqlPartitioningRepository.java
  66. 14
      dao/src/main/resources/sql/schema-entities.sql
  67. 2
      rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java
  68. 1
      rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/slack/SlackService.java
  69. 6
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbNotificationNode.java
  70. 7
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbNotificationNodeConfiguration.java
  71. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbSlackNode.java
  72. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbSlackNodeConfiguration.java

14
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)
);

5
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;

5
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) {

12
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> T checkNotNull(T reference) throws ThingsboardException {

3
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");
}

7
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);
}
}

20
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<User> recipients = notificationTargetService.findRecipientsForNotificationTargetConfig(user.getTenantId(), notificationTarget.getConfiguration(), null);
for (User recipient : recipients.getData()) {
accessControlService.checkPermission(user, Resource.USER, Operation.READ, recipient.getId(), recipient);
}
PageDataIterable<User> 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<User> recipients = notificationTargetService.findRecipientsForNotificationTargetConfig(user.getTenantId(), notificationTarget.getConfiguration(), pageLink);
PageData<User> 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);
}
}

3
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);
}

29
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]
}
}

5
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
}
}

3
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();
}

80
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<ListenableFuture<Void>> 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<ListenableFuture<Void>> 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<User> recipients = recipientsBatch.getData();
for (User recipient : recipients) {
ListenableFuture<Void> 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<User> recipients = recipientsBatch.getData();
for (User recipient : recipients) {
ListenableFuture<Void> 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<Void> 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

53
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<Void> onAlarmCreatedOrUpdated(TenantId tenantId, Alarm alarm) {
@ -67,23 +73,12 @@ public class DefaultNotificationRuleProcessingService implements NotificationRul
private ListenableFuture<Void> 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<Void> 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<NotificationRequest> 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<NotificationRequestId> scheduledForRule = notificationRequestService.findNotificationRequestsIdsByStatusAndRuleId(tenantId, NotificationRequestStatus.SCHEDULED, notificationRuleId);
for (NotificationRequestId notificationRequestId : scheduledForRule) {
notificationManager.deleteNotificationRequest(tenantId, notificationRequestId);
}
}
}

59
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<NotificationRequestId, ScheduledFuture<?>> scheduledNotificationRequests = new ConcurrentHashMap<>();
private final Map<NotificationRequestId, ScheduledRequestMetadata> 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<NotificationRequestId> 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;
}
}

34
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<NotificationDeliveryMethod, DeliveryMethodNotificationTemplate> templates;
@Getter
private NotificationRequestStats stats;
private final NotificationRequestStats stats;
@Builder
public NotificationProcessingContext(TenantId tenantId, NotificationSettings settings, NotificationRequest request, Map<String, String> 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 extends DeliveryMethodNotificationTemplate> T getTemplate(NotificationDeliveryMethod deliveryMethod) {
@ -90,4 +94,8 @@ public class NotificationProcessingContext {
return templateContext;
}
public CustomerId getOriginatorCustomerId() {
return originatorEntity != null ? originatorEntity.getCustomerId() : null;
}
}

3
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<Void> onAlarmDeleted(TenantId tenantId, Alarm alarm);
ListenableFuture<Void> onNotificationRuleDeleted(TenantId tenantId, NotificationRuleId ruleId);
}

2
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);
}

2
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;

15
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<ToCore
OtaPackageStateService firmwareStateService,
GitVersionControlQueueService vcQueueService,
PartitionService partitionService,
ApplicationEventPublisher eventPublisher,
Optional<JwtSettingsService> 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<ToCore
private void forwardToNotificationSchedulerService(TransportProtos.NotificationSchedulerServiceMsg msg, TbCallback callback) {
TenantId tenantId = TenantId.fromUUID(new UUID(msg.getTenantIdMSB(), msg.getTenantIdLSB()));
NotificationRequestId notificationRequestId = new NotificationRequestId(new UUID(msg.getRequestIdMSB(), msg.getRequestIdLSB()));
boolean deleted = msg.getDeleted();
try {
if (!deleted) {
notificationSchedulerService.scheduleNotificationRequest(tenantId, notificationRequestId, msg.getTs());
} else {
notificationSchedulerService.onNotificationRequestDeleted(tenantId, notificationRequestId);
}
notificationSchedulerService.scheduleNotificationRequest(tenantId, notificationRequestId, msg.getTs());
callback.onSuccess();
} catch (Exception e) {
callback.onFailure(new RuntimeException("Failed to scheduler notification request", e));

6
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java

@ -18,6 +18,7 @@ package org.thingsboard.server.service.queue;
import com.google.protobuf.ProtocolStringList;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.ApplicationEventPublisher;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Service;
import org.thingsboard.common.util.ThingsBoardThreadFactory;
@ -126,8 +127,9 @@ public class DefaultTbRuleEngineConsumerService extends AbstractConsumerService<
TbAssetProfileCache assetProfileCache,
TbTenantProfileCache tenantProfileCache,
TbApiUsageStateService apiUsageStateService,
PartitionService partitionService, TbServiceInfoProvider serviceInfoProvider, QueueService queueService) {
super(actorContext, encodingService, tenantProfileCache, deviceProfileCache, assetProfileCache, apiUsageStateService, partitionService, tbRuleEngineQueueFactory.createToRuleEngineNotificationsMsgConsumer(), Optional.empty());
PartitionService partitionService, ApplicationEventPublisher eventPublisher,
TbServiceInfoProvider serviceInfoProvider, QueueService queueService) {
super(actorContext, encodingService, tenantProfileCache, deviceProfileCache, assetProfileCache, apiUsageStateService, partitionService, eventPublisher, tbRuleEngineQueueFactory.createToRuleEngineNotificationsMsgConsumer(), Optional.empty());
this.statisticsService = statisticsService;
this.tbRuleEngineQueueFactory = tbRuleEngineQueueFactory;
this.submitStrategyFactory = submitStrategyFactory;

7
application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java

@ -18,6 +18,7 @@ package org.thingsboard.server.service.queue.processing;
import com.google.protobuf.ByteString;
import lombok.extern.slf4j.Slf4j;
import org.springframework.boot.context.event.ApplicationReadyEvent;
import org.springframework.context.ApplicationEventPublisher;
import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.actors.ActorSystemContext;
import org.thingsboard.server.common.data.EntityType;
@ -75,6 +76,7 @@ public abstract class AbstractConsumerService<N extends com.google.protobuf.Gene
protected final TbAssetProfileCache assetProfileCache;
protected final TbApiUsageStateService apiUsageStateService;
protected final PartitionService partitionService;
protected final ApplicationEventPublisher eventPublisher;
protected final TbQueueConsumer<TbProtoQueueMsg<N>> nfConsumer;
protected final Optional<JwtSettingsService> jwtSettingsService;
@ -83,7 +85,8 @@ public abstract class AbstractConsumerService<N extends com.google.protobuf.Gene
public AbstractConsumerService(ActorSystemContext actorContext, DataDecodingEncodingService encodingService,
TbTenantProfileCache tenantProfileCache, TbDeviceProfileCache deviceProfileCache,
TbAssetProfileCache assetProfileCache, TbApiUsageStateService apiUsageStateService,
PartitionService partitionService, TbQueueConsumer<TbProtoQueueMsg<N>> nfConsumer, Optional<JwtSettingsService> jwtSettingsService) {
PartitionService partitionService, ApplicationEventPublisher eventPublisher,
TbQueueConsumer<TbProtoQueueMsg<N>> nfConsumer, Optional<JwtSettingsService> jwtSettingsService) {
this.actorContext = actorContext;
this.encodingService = encodingService;
this.tenantProfileCache = tenantProfileCache;
@ -91,6 +94,7 @@ public abstract class AbstractConsumerService<N extends com.google.protobuf.Gene
this.assetProfileCache = assetProfileCache;
this.apiUsageStateService = apiUsageStateService;
this.partitionService = partitionService;
this.eventPublisher = eventPublisher;
this.nfConsumer = nfConsumer;
this.jwtSettingsService = jwtSettingsService;
}
@ -205,6 +209,7 @@ public abstract class AbstractConsumerService<N extends com.google.protobuf.Gene
apiUsageStateService.onCustomerDelete((CustomerId) componentLifecycleMsg.getEntityId());
}
}
eventPublisher.publishEvent(componentLifecycleMsg);
}
log.trace("[{}] Forwarding message to App Actor {}", id, actorMsg);
actorContext.tellWithHighPriority(actorMsg);

2
application/src/main/java/org/thingsboard/server/service/slack/DefaultSlackService.java

@ -30,7 +30,7 @@ import com.slack.api.model.ConversationType;
import lombok.RequiredArgsConstructor;
import org.apache.commons.lang3.StringUtils;
import org.springframework.stereotype.Service;
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.id.TenantId;
import org.thingsboard.server.common.data.notification.NotificationDeliveryMethod;

20
application/src/main/java/org/thingsboard/server/service/ttl/NotificationsCleanUpService.java

@ -20,6 +20,8 @@ import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.notification.NotificationRequestConfig;
import org.thingsboard.server.dao.notification.NotificationRequestDao;
import org.thingsboard.server.dao.sqlts.insert.sql.SqlPartitioningRepository;
import org.thingsboard.server.queue.discovery.PartitionService;
@ -33,15 +35,18 @@ import static org.thingsboard.server.dao.model.ModelConstants.NOTIFICATION_TABLE
public class NotificationsCleanUpService extends AbstractCleanUpService {
private final SqlPartitioningRepository partitioningRepository;
private final NotificationRequestDao notificationRequestDao;
@Value("${sql.ttl.notifications.ttl:2592000}")
private long ttlInSec;
@Value("${sql.notifications.partition_size:168}")
private int partitionSizeInHours;
public NotificationsCleanUpService(PartitionService partitionService, SqlPartitioningRepository partitioningRepository) {
public NotificationsCleanUpService(PartitionService partitionService, SqlPartitioningRepository partitioningRepository,
NotificationRequestDao notificationRequestDao) {
super(partitionService);
this.partitioningRepository = partitioningRepository;
this.notificationRequestDao = notificationRequestDao;
}
@Scheduled(initialDelayString = "#{T(org.apache.commons.lang3.RandomUtils).nextLong(0, ${sql.ttl.notifications.checking_interval_ms:86400000})}",
@ -49,10 +54,17 @@ public class NotificationsCleanUpService extends AbstractCleanUpService {
public void cleanUp() {
long expTime = System.currentTimeMillis() - TimeUnit.SECONDS.toMillis(ttlInSec);
long partitionDurationMs = TimeUnit.HOURS.toMillis(partitionSizeInHours);
if (isSystemTenantPartitionMine()) {
partitioningRepository.dropPartitionsBefore(NOTIFICATION_TABLE_NAME, expTime, partitionDurationMs);
} else {
if (!isSystemTenantPartitionMine()) {
partitioningRepository.cleanupPartitionsCache(NOTIFICATION_TABLE_NAME, expTime, partitionDurationMs);
return;
}
long lastRemovedNotificationTs = partitioningRepository.dropPartitionsBefore(NOTIFICATION_TABLE_NAME, expTime, partitionDurationMs);
if (lastRemovedNotificationTs > 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);
}
}

27
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);

23
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);

168
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<User> 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<User> getRecipients(NotificationTarget notificationTarget) throws Exception {
return doPostWithTypedResponse("/api/notification/target/recipients?page=0&pageSize=100", notificationTarget, new TypeReference<PageData<User>>() {}).getData();
}
}

90
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("<script/>");
validationError = saveAndGetError(notificationTemplate, status().isBadRequest());
assertThat(validationError)
.doesNotContain("defaultTextTemplate must be specified")
.contains("subject must not be")
.contains("name is malformed");
config.setDefaultTextTemplate(null);
validationError = saveAndGetError(notificationTemplate, status().isBadRequest());
assertThat(validationError)
.contains("defaultTextTemplate must be specified");
}
private String saveAndGetError(NotificationTemplate notificationTemplate, ResultMatcher statusMatcher) throws Exception {
return getErrorMessage(save(notificationTemplate, statusMatcher));
}
private ResultActions save(NotificationTemplate notificationTemplate, ResultMatcher statusMatcher) throws Exception {
return doPost("/api/notification/template", notificationTemplate)
.andExpect(statusMatcher);
}
}

86
application/src/test/java/org/thingsboard/server/service/notification/NotificationsClient.java

@ -1,86 +0,0 @@
/**
* 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.google.common.base.Strings;
import org.apache.commons.lang3.StringUtils;
import org.thingsboard.rest.client.RestClient;
import org.thingsboard.server.common.data.notification.AlarmOriginatedNotificationInfo;
import org.thingsboard.server.common.data.notification.Notification;
import org.thingsboard.server.common.data.notification.NotificationInfo;
import org.thingsboard.server.common.data.notification.NotificationOriginatorType;
import java.text.SimpleDateFormat;
import java.util.Date;
import java.util.List;
import java.util.Scanner;
public class NotificationsClient extends NotificationApiWsClient {
private NotificationsClient(String wsUrl, String token) throws Exception {
super(wsUrl, token);
}
public static NotificationsClient newInstance(String username, String password) throws Exception {
RestClient restClient = new RestClient("http://localhost:8080");
restClient.login(username, password);
NotificationsClient client = new NotificationsClient("ws://localhost:8080", restClient.getToken());
client.connectBlocking();
return client;
}
@Override
public void onMessage(String s) {
super.onMessage(s);
// printNotificationsCount();
printNotifications();
}
public void printNotifications() {
System.out.println(StringUtils.repeat(System.lineSeparator(), 20));
List<Notification> notifications = getNotifications();
System.out.printf(" %s NEW MESSAGE%s\n\n", getUnreadCount(), notifications.size() > 1 ? "S" : "");
notifications.forEach(notification -> {
String notificationInfoStr = "";
if (notification.getOriginatorType() == NotificationOriginatorType.ALARM) {
AlarmOriginatedNotificationInfo info = (AlarmOriginatedNotificationInfo) notification.getInfo();
notificationInfoStr = String.format("Alarm of type %s - %s severity - status: %s",
info.getAlarmType(), info.getAlarmSeverity(), info.getAlarmStatus());
} else if (notification.getInfo() != null) {
notificationInfoStr = Strings.nullToEmpty(notification.getInfo().getDescription());
}
SimpleDateFormat format = new SimpleDateFormat("dd.MM.yyyy HH:mm:ss");
String time = format.format(new Date(notification.getCreatedTime()));
// System.out.printf("[%s] %-19s | %-30s | (%s)\n", time, notification.getReason(), notification.getText(), notificationInfoStr);
});
System.out.println(StringUtils.repeat(System.lineSeparator(), 5));
}
public void printNotificationsCount() {
System.out.println();
System.out.println();
System.out.println();
int unreadCount = getUnreadCount();
System.out.printf("\r\r%s NEW MESSAGE%s", unreadCount, unreadCount > 1 ? "S" : "");
}
public static void main(String[] args) throws Exception {
NotificationsClient client = NotificationsClient.newInstance("tenant@thingsboard.org", "tenant");
client.subscribeForUnreadNotifications(5);
// client.subscribeForUnreadNotificationsCount();
new Scanner(System.in).nextLine();
}
}

1
common/cluster-api/src/main/proto/queue.proto

@ -1042,5 +1042,4 @@ message NotificationSchedulerServiceMsg {
int64 requestIdMSB = 3;
int64 requestIdLSB = 4;
int64 ts = 5;
bool deleted = 6;
}

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

@ -21,6 +21,7 @@ import org.thingsboard.server.common.data.id.NotificationRuleId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.notification.NotificationRequest;
import org.thingsboard.server.common.data.notification.NotificationRequestStats;
import org.thingsboard.server.common.data.notification.NotificationRequestStatus;
import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageLink;
@ -34,6 +35,8 @@ public interface NotificationRequestService {
PageData<NotificationRequest> findNotificationRequestsByTenantId(TenantId tenantId, PageLink pageLink);
List<NotificationRequestId> findNotificationRequestsIdsByStatusAndRuleId(TenantId tenantId, NotificationRequestStatus requestStatus, NotificationRuleId ruleId);
List<NotificationRequest> findNotificationRequestsByRuleIdAndOriginatorEntityId(TenantId tenantId, NotificationRuleId ruleId, EntityId originatorEntityId);
void deleteNotificationRequestById(TenantId tenantId, NotificationRequestId id);

4
common/dao-api/src/main/java/org/thingsboard/server/dao/notification/NotificationRuleService.java

@ -25,10 +25,10 @@ public interface NotificationRuleService {
NotificationRule saveNotificationRule(TenantId tenantId, NotificationRule notificationRule);
NotificationRule findNotificationRuleById(TenantId tenantId, NotificationRuleId notificationRuleId);
NotificationRule findNotificationRuleById(TenantId tenantId, NotificationRuleId id);
PageData<NotificationRule> findNotificationRulesByTenantId(TenantId tenantId, PageLink pageLink);
void deleteNotificationRule(TenantId tenantId, NotificationRuleId notificationRuleId);
void deleteNotificationRuleById(TenantId tenantId, NotificationRuleId id);
}

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

@ -16,6 +16,7 @@
package org.thingsboard.server.dao.notification;
import org.thingsboard.server.common.data.User;
import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.NotificationTargetId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.notification.targets.NotificationTarget;
@ -31,10 +32,10 @@ public interface NotificationTargetService {
PageData<NotificationTarget> findNotificationTargetsByTenantId(TenantId tenantId, PageLink pageLink);
PageData<User> findRecipientsForNotificationTarget(TenantId tenantId, NotificationTargetId notificationTargetId, PageLink pageLink);
PageData<User> findRecipientsForNotificationTarget(TenantId tenantId, CustomerId customerId, NotificationTargetId notificationTargetId, PageLink pageLink);
PageData<User> findRecipientsForNotificationTargetConfig(TenantId tenantId, NotificationTargetConfig targetConfig, PageLink pageLink);
PageData<User> findRecipientsForNotificationTargetConfig(TenantId tenantId, CustomerId customerId, NotificationTargetConfig targetConfig, PageLink pageLink);
void deleteNotificationTarget(TenantId tenantId, NotificationTargetId notificationTargetId);
void deleteNotificationTargetById(TenantId tenantId, NotificationTargetId id);
}

14
common/data/src/main/java/org/thingsboard/server/common/data/notification/NotificationRequest.java

@ -23,6 +23,7 @@ import lombok.EqualsAndHashCode;
import lombok.NoArgsConstructor;
import org.apache.commons.lang3.StringUtils;
import org.thingsboard.server.common.data.BaseData;
import org.thingsboard.server.common.data.HasCustomerId;
import org.thingsboard.server.common.data.HasName;
import org.thingsboard.server.common.data.HasTenantId;
import org.thingsboard.server.common.data.id.EntityId;
@ -46,8 +47,8 @@ import java.util.Map;
public class NotificationRequest extends BaseData<NotificationRequestId> implements HasTenantId, HasName {
private TenantId tenantId;
@NotNull
private NotificationTargetId targetId;
@NotEmpty
private List<NotificationTargetId> targets;
@NotNull
private NotificationTemplateId templateId;
@ -68,11 +69,18 @@ public class NotificationRequest extends BaseData<NotificationRequestId> impleme
private NotificationRequestStats stats;
@JsonIgnore
private transient Map<String, String> templateContext;
@JsonIgnore
private transient HasCustomerId originatorEntity;
public void copyContext(NotificationRequest other) {
this.templateContext = other.getTemplateContext();
this.originatorEntity = other.getOriginatorEntity();
}
@JsonIgnore
@Override
public String getName() {
return "To target " + targetId + " via " + StringUtils.join(deliveryMethods, ", ");
return "To targets " + targets + " via " + StringUtils.join(deliveryMethods, ", ");
}
}

4
common/data/src/main/java/org/thingsboard/server/common/data/notification/NotificationRequestConfig.java

@ -22,7 +22,9 @@ import javax.validation.constraints.Max;
@Data
public class NotificationRequestConfig {
@Max(value = 604800, message = "cannot be longer than 1 week")
@Max(value = MAX_SENDING_DELAY, message = "cannot be longer than 1 week")
private int sendingDelayInSec;
public static final int MAX_SENDING_DELAY = 604800;
}

15
common/data/src/main/java/org/thingsboard/server/common/data/notification/NotificationRequestStats.java

@ -16,11 +16,14 @@
package org.thingsboard.server.common.data.notification;
import com.fasterxml.jackson.annotation.JsonCreator;
import com.fasterxml.jackson.annotation.JsonIgnore;
import com.fasterxml.jackson.annotation.JsonProperty;
import lombok.Data;
import org.thingsboard.server.common.data.User;
import org.thingsboard.server.common.data.id.UserId;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicInteger;
@ -29,10 +32,13 @@ public class NotificationRequestStats {
private final Map<NotificationDeliveryMethod, AtomicInteger> sent;
private final Map<NotificationDeliveryMethod, Map<String, String>> errors;
@JsonIgnore
private final Map<NotificationDeliveryMethod, Set<UserId>> processedRecipients;
public NotificationRequestStats() {
this.sent = new ConcurrentHashMap<>();
this.errors = new ConcurrentHashMap<>();
this.processedRecipients = new ConcurrentHashMap<>();
}
@JsonCreator
@ -40,10 +46,12 @@ public class NotificationRequestStats {
@JsonProperty("errors") Map<NotificationDeliveryMethod, Map<String, String>> errors) {
this.sent = sent;
this.errors = errors;
this.processedRecipients = null;
}
public void reportSent(NotificationDeliveryMethod deliveryMethod) {
public void reportSent(NotificationDeliveryMethod deliveryMethod, User recipient) {
sent.computeIfAbsent(deliveryMethod, k -> new AtomicInteger()).incrementAndGet();
processedRecipients.computeIfAbsent(deliveryMethod, k -> ConcurrentHashMap.newKeySet()).add(recipient.getId());
}
public void reportError(NotificationDeliveryMethod deliveryMethod, User recipient, Throwable error) {
@ -58,4 +66,9 @@ public class NotificationRequestStats {
return sent.containsKey(deliveryMethod) || errors.containsKey(deliveryMethod);
}
public boolean contains(NotificationDeliveryMethod deliveryMethod, UserId recipientId) {
Set<UserId> processedRecipients = this.processedRecipients.get(deliveryMethod);
return processedRecipients != null && processedRecipients.contains(recipientId);
}
}

5
common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/NonConfirmedNotificationEscalation.java

@ -17,14 +17,15 @@ package org.thingsboard.server.common.data.notification.rule;
import lombok.Data;
import org.thingsboard.server.common.data.id.NotificationTargetId;
import org.thingsboard.server.common.data.notification.NotificationRequestConfig;
import javax.validation.constraints.Min;
import javax.validation.constraints.Max;
import javax.validation.constraints.NotNull;
@Data
public class NonConfirmedNotificationEscalation {
@Min(1)
@Max(NotificationRequestConfig.MAX_SENDING_DELAY)
private int delayInSec;
@NotNull
private NotificationTargetId notificationTargetId;

16
common/data/src/main/java/org/thingsboard/server/common/data/notification/targets/CustomerUsersNotificationTargetConfig.java

@ -15,18 +15,32 @@
*/
package org.thingsboard.server.common.data.notification.targets;
import com.fasterxml.jackson.annotation.JsonIgnore;
import lombok.Data;
import org.thingsboard.server.common.data.id.EntityId;
import javax.validation.constraints.AssertTrue;
import java.util.UUID;
@Data
public class CustomerUsersNotificationTargetConfig implements NotificationTargetConfig {
private UUID customerId;
private UUID customerId; // might not be set if using with notification rule
private boolean getCustomerIdFromOriginatorEntity; // e.g. from alarm
@Override
public NotificationTargetConfigType getType() {
return NotificationTargetConfigType.CUSTOMER_USERS;
}
@AssertTrue(message = "customerId is required")
@JsonIgnore
public boolean isValid() {
if (!getCustomerIdFromOriginatorEntity) {
return customerId != null && !customerId.equals(EntityId.NULL_UUID);
} else {
return true;
}
}
}

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

@ -22,6 +22,7 @@ import org.thingsboard.server.common.data.HasName;
import org.thingsboard.server.common.data.HasTenantId;
import org.thingsboard.server.common.data.id.NotificationTargetId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.validation.NoXss;
import javax.validation.Valid;
import javax.validation.constraints.NotBlank;
@ -31,9 +32,9 @@ import javax.validation.constraints.NotNull;
@EqualsAndHashCode(callSuper = true)
public class NotificationTarget extends BaseData<NotificationTargetId> implements HasTenantId, HasName {
@NotNull
private TenantId tenantId;
@NotBlank
@NoXss
private String name;
@NotNull
@Valid

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

@ -21,6 +21,8 @@ import com.fasterxml.jackson.annotation.JsonTypeInfo;
import lombok.Data;
import org.thingsboard.server.common.data.notification.NotificationDeliveryMethod;
import javax.validation.constraints.NotNull;
@JsonIgnoreProperties(ignoreUnknown = true)
@JsonTypeInfo(use = JsonTypeInfo.Id.NAME, property = "method", visible = true, include = JsonTypeInfo.As.EXISTING_PROPERTY, defaultImpl = DeliveryMethodNotificationTemplate.class)
@JsonSubTypes({
@ -31,6 +33,7 @@ import org.thingsboard.server.common.data.notification.NotificationDeliveryMetho
public class DeliveryMethodNotificationTemplate {
private String body;
@NotNull
private NotificationDeliveryMethod method;
}

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

@ -18,10 +18,13 @@ package org.thingsboard.server.common.data.notification.template;
import lombok.Data;
import lombok.EqualsAndHashCode;
import javax.validation.constraints.NotBlank;
@Data
@EqualsAndHashCode(callSuper = true)
public class EmailDeliveryMethodNotificationTemplate extends DeliveryMethodNotificationTemplate {
@NotBlank
private String subject;
}

1
common/data/src/main/java/org/thingsboard/server/common/data/notification/template/NotificationTemplate.java

@ -39,6 +39,7 @@ public class NotificationTemplate extends BaseData<NotificationTemplateId> imple
@NotNull
private String notificationType;
@Valid
@NotNull
private NotificationTemplateConfig configuration;
}

19
common/data/src/main/java/org/thingsboard/server/common/data/notification/template/NotificationTemplateConfig.java

@ -15,15 +15,32 @@
*/
package org.thingsboard.server.common.data.notification.template;
import com.fasterxml.jackson.annotation.JsonIgnore;
import lombok.Data;
import org.apache.commons.lang3.StringUtils;
import org.thingsboard.server.common.data.notification.NotificationDeliveryMethod;
import javax.validation.Valid;
import javax.validation.constraints.AssertTrue;
import javax.validation.constraints.NotEmpty;
import java.util.Map;
@Data
public class NotificationTemplateConfig {
private DeliveryMethodNotificationTemplate defaultTemplate;
private String defaultTextTemplate;
@Valid
@NotEmpty
private Map<NotificationDeliveryMethod, DeliveryMethodNotificationTemplate> templates;
@JsonIgnore
@AssertTrue(message = "defaultTextTemplate must be specified if one absent for delivery method")
public boolean isValid() {
if (templates.values().stream().anyMatch(template -> StringUtils.isEmpty(template.getBody()))) {
return StringUtils.isNotEmpty(defaultTextTemplate);
} else {
return true;
}
}
}

2
rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/slack/SlackConversation.java → common/data/src/main/java/org/thingsboard/server/common/data/notification/template/SlackConversation.java

@ -13,7 +13,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.rule.engine.api.slack;
package org.thingsboard.server.common.data.notification.template;
import lombok.Data;

2
common/data/src/main/java/org/thingsboard/server/common/data/notification/template/SlackDeliveryMethodNotificationTemplate.java

@ -22,6 +22,8 @@ import lombok.EqualsAndHashCode;
@EqualsAndHashCode(callSuper = true)
public class SlackDeliveryMethodNotificationTemplate extends DeliveryMethodNotificationTemplate {
private SlackConversation.Type conversationType;
// add Conversation type!
private String conversationId; // not required, set from user's name if not set
}

44
dao/src/main/java/org/thingsboard/server/dao/model/BaseSqlEntity.java

@ -15,17 +15,23 @@
*/
package org.thingsboard.server.dao.model;
import com.fasterxml.jackson.core.type.TypeReference;
import com.fasterxml.jackson.databind.JsonNode;
import lombok.Data;
import org.apache.commons.lang3.StringUtils;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.id.UUIDBased;
import javax.persistence.Column;
import javax.persistence.Id;
import javax.persistence.MappedSuperclass;
import java.util.Arrays;
import java.util.Collections;
import java.util.List;
import java.util.UUID;
import java.util.function.Function;
import java.util.stream.Collectors;
/**
* Created by ashvayka on 13.07.17.
@ -70,7 +76,15 @@ public abstract class BaseSqlEntity<D> implements BaseEntity<D> {
}
}
protected static <I> I createId(UUID uuid, Function<UUID, I> creator) {
protected static UUID getTenantUuid(TenantId tenantId) {
if (tenantId != null && !tenantId.isNullUid()) {
return tenantId.getId();
} else {
return null;
}
}
protected static <I> I getEntityId(UUID uuid, Function<UUID, I> creator) {
if (uuid != null) {
return creator.apply(uuid);
} else {
@ -78,6 +92,14 @@ public abstract class BaseSqlEntity<D> implements BaseEntity<D> {
}
}
protected static TenantId getTenantId(UUID uuid) {
if (uuid != null && !uuid.equals(EntityId.NULL_UUID)) {
return TenantId.fromUUID(uuid);
} else {
return TenantId.SYS_TENANT_ID;
}
}
protected JsonNode toJson(Object value) {
if (value != null) {
return JacksonUtil.valueToTree(value);
@ -90,4 +112,22 @@ public abstract class BaseSqlEntity<D> implements BaseEntity<D> {
return JacksonUtil.convertValue(json, type);
}
protected String listToString(List<?> list) {
if (list != null) {
return StringUtils.join(list, ',');
} else {
return "";
}
}
protected <E> List<E> listFromString(String string, Function<String, E> mappingFunction) {
if (string != null) {
return Arrays.stream(StringUtils.split(string, ','))
.filter(StringUtils::isNotBlank)
.map(mappingFunction).collect(Collectors.toList());
} else {
return Collections.emptyList();
}
}
}

2
dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java

@ -662,7 +662,7 @@ public class ModelConstants {
public static final String NOTIFICATION_STATUS_PROPERTY = "status";
public static final String NOTIFICATION_REQUEST_TABLE_NAME = "notification_request";
public static final String NOTIFICATION_REQUEST_TARGET_ID_PROPERTY = "target_id";
public static final String NOTIFICATION_REQUEST_TARGETS_PROPERTY = "targets";
public static final String NOTIFICATION_REQUEST_TEMPLATE_ID_PROPERTY = "template_id";
public static final String NOTIFICATION_REQUEST_DELIVERY_METHODS_PROPERTY = "delivery_methods";
public static final String NOTIFICATION_REQUEST_INFO_PROPERTY = "info";

2
dao/src/main/java/org/thingsboard/server/dao/model/sql/AbstractAlarmEntity.java

@ -207,7 +207,7 @@ public abstract class AbstractAlarmEntity<T extends Alarm> extends BaseSqlEntity
} else {
alarm.setPropagateRelationTypes(Collections.emptyList());
}
alarm.setNotificationRuleId(createId(notificationRuleId, NotificationRuleId::new));
alarm.setNotificationRuleId(getEntityId(notificationRuleId, NotificationRuleId::new));
return alarm;
}
}

5
dao/src/main/java/org/thingsboard/server/dao/model/sql/NotificationEntity.java

@ -21,7 +21,6 @@ import lombok.EqualsAndHashCode;
import org.hibernate.annotations.Formula;
import org.hibernate.annotations.Type;
import org.hibernate.annotations.TypeDef;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.common.data.id.NotificationId;
import org.thingsboard.server.common.data.id.NotificationRequestId;
import org.thingsboard.server.common.data.id.UserId;
@ -90,8 +89,8 @@ public class NotificationEntity extends BaseSqlEntity<Notification> {
Notification notification = new Notification();
notification.setId(new NotificationId(id));
notification.setCreatedTime(createdTime);
notification.setRequestId(createId(requestId, NotificationRequestId::new));
notification.setRecipientId(createId(recipientId, UserId::new));
notification.setRequestId(getEntityId(requestId, NotificationRequestId::new));
notification.setRecipientId(getEntityId(recipientId, UserId::new));
notification.setText(type);
notification.setText(text);
notification.setInfo(fromJson(info, NotificationInfo.class));

27
dao/src/main/java/org/thingsboard/server/dao/model/sql/NotificationRequestEntity.java

@ -18,7 +18,6 @@ package org.thingsboard.server.dao.model.sql;
import com.fasterxml.jackson.databind.JsonNode;
import lombok.Data;
import lombok.EqualsAndHashCode;
import org.apache.commons.lang3.StringUtils;
import org.hibernate.annotations.Type;
import org.hibernate.annotations.TypeDef;
import org.thingsboard.server.common.data.EntityType;
@ -27,7 +26,6 @@ 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.NotificationTemplateId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.notification.NotificationDeliveryMethod;
import org.thingsboard.server.common.data.notification.NotificationInfo;
import org.thingsboard.server.common.data.notification.NotificationOriginatorType;
@ -44,9 +42,7 @@ import javax.persistence.Entity;
import javax.persistence.EnumType;
import javax.persistence.Enumerated;
import javax.persistence.Table;
import java.util.Arrays;
import java.util.UUID;
import java.util.stream.Collectors;
@Data
@EqualsAndHashCode(callSuper = true)
@ -58,8 +54,8 @@ public class NotificationRequestEntity extends BaseSqlEntity<NotificationRequest
@Column(name = ModelConstants.TENANT_ID_PROPERTY, nullable = false)
private UUID tenantId;
@Column(name = ModelConstants.NOTIFICATION_REQUEST_TARGET_ID_PROPERTY, nullable = false)
private UUID targetId;
@Column(name = ModelConstants.NOTIFICATION_REQUEST_TARGETS_PROPERTY, nullable = false)
private String targets;
@Column(name = ModelConstants.NOTIFICATION_REQUEST_TEMPLATE_ID_PROPERTY, nullable = false)
private UUID templateId;
@ -102,11 +98,11 @@ public class NotificationRequestEntity extends BaseSqlEntity<NotificationRequest
public NotificationRequestEntity(NotificationRequest notificationRequest) {
setId(notificationRequest.getUuidId());
setCreatedTime(notificationRequest.getCreatedTime());
setTenantId(getUuid(notificationRequest.getTenantId()));
setTargetId(getUuid(notificationRequest.getTargetId()));
setTenantId(getTenantUuid(notificationRequest.getTenantId()));
setTargets(listToString(notificationRequest.getTargets()));
setTemplateId(getUuid(notificationRequest.getTemplateId()));
setInfo(toJson(notificationRequest.getInfo()));
setDeliveryMethods(StringUtils.join(notificationRequest.getDeliveryMethods(), ','));
setDeliveryMethods(listToString(notificationRequest.getDeliveryMethods()));
setAdditionalConfig(toJson(notificationRequest.getAdditionalConfig()));
setOriginatorType(notificationRequest.getOriginatorType());
if (notificationRequest.getOriginatorEntityId() != null) {
@ -123,20 +119,17 @@ public class NotificationRequestEntity extends BaseSqlEntity<NotificationRequest
NotificationRequest notificationRequest = new NotificationRequest();
notificationRequest.setId(new NotificationRequestId(id));
notificationRequest.setCreatedTime(createdTime);
notificationRequest.setTenantId(createId(tenantId, TenantId::new));
notificationRequest.setTargetId(createId(targetId, NotificationTargetId::new));
notificationRequest.setTemplateId(createId(templateId, NotificationTemplateId::new));
notificationRequest.setTenantId(getTenantId(tenantId));
notificationRequest.setTargets(listFromString(targets, uuid -> new NotificationTargetId(UUID.fromString(uuid))));
notificationRequest.setTemplateId(getEntityId(templateId, NotificationTemplateId::new));
notificationRequest.setInfo(fromJson(info, NotificationInfo.class));
if (deliveryMethods != null) {
notificationRequest.setDeliveryMethods(Arrays.stream(StringUtils.split(deliveryMethods, ','))
.filter(StringUtils::isNotBlank).map(NotificationDeliveryMethod::valueOf).collect(Collectors.toList()));
}
notificationRequest.setDeliveryMethods(listFromString(deliveryMethods, NotificationDeliveryMethod::valueOf));
notificationRequest.setAdditionalConfig(fromJson(additionalConfig, NotificationRequestConfig.class));
notificationRequest.setOriginatorType(originatorType);
if (originatorEntityId != null) {
notificationRequest.setOriginatorEntityId(EntityIdFactory.getByTypeAndUuid(originatorEntityType, originatorEntityId));
}
notificationRequest.setRuleId(createId(ruleId, NotificationRuleId::new));
notificationRequest.setRuleId(getEntityId(ruleId, NotificationRuleId::new));
notificationRequest.setStatus(status);
notificationRequest.setStats(fromJson(stats, NotificationRequestStats.class));
return notificationRequest;

17
dao/src/main/java/org/thingsboard/server/dao/model/sql/NotificationRuleEntity.java

@ -18,12 +18,10 @@ package org.thingsboard.server.dao.model.sql;
import com.fasterxml.jackson.databind.JsonNode;
import lombok.Data;
import lombok.EqualsAndHashCode;
import org.apache.commons.lang3.StringUtils;
import org.hibernate.annotations.Type;
import org.hibernate.annotations.TypeDef;
import org.thingsboard.server.common.data.id.NotificationRuleId;
import org.thingsboard.server.common.data.id.NotificationTemplateId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.notification.NotificationDeliveryMethod;
import org.thingsboard.server.common.data.notification.rule.NotificationRule;
import org.thingsboard.server.common.data.notification.rule.NotificationRuleConfig;
@ -34,9 +32,7 @@ import org.thingsboard.server.dao.util.mapping.JsonStringType;
import javax.persistence.Column;
import javax.persistence.Entity;
import javax.persistence.Table;
import java.util.Arrays;
import java.util.UUID;
import java.util.stream.Collectors;
@Data @EqualsAndHashCode(callSuper = true)
@Entity
@ -65,10 +61,10 @@ public class NotificationRuleEntity extends BaseSqlEntity<NotificationRule> {
public NotificationRuleEntity(NotificationRule notificationRule) {
setId(notificationRule.getUuidId());
setCreatedTime(notificationRule.getCreatedTime());
setTenantId(getUuid(notificationRule.getTenantId()));
setTenantId(getTenantUuid(notificationRule.getTenantId()));
setName(notificationRule.getName());
setTemplateId(getUuid(notificationRule.getTemplateId()));
setDeliveryMethods(StringUtils.join(notificationRule.getDeliveryMethods(), ','));
setDeliveryMethods(listToString(notificationRule.getDeliveryMethods()));
setConfiguration(toJson(notificationRule.getConfiguration()));
}
@ -77,13 +73,10 @@ public class NotificationRuleEntity extends BaseSqlEntity<NotificationRule> {
NotificationRule notificationRule = new NotificationRule();
notificationRule.setId(new NotificationRuleId(id));
notificationRule.setCreatedTime(createdTime);
notificationRule.setTenantId(createId(tenantId, TenantId::fromUUID));
notificationRule.setTenantId(getTenantId(tenantId));
notificationRule.setName(name);
notificationRule.setTemplateId(createId(templateId, NotificationTemplateId::new));
if (deliveryMethods != null) {
notificationRule.setDeliveryMethods(Arrays.stream(StringUtils.split(deliveryMethods, ','))
.filter(StringUtils::isNotBlank).map(NotificationDeliveryMethod::valueOf).collect(Collectors.toList()));
}
notificationRule.setTemplateId(getEntityId(templateId, NotificationTemplateId::new));
notificationRule.setDeliveryMethods(listFromString(deliveryMethods, NotificationDeliveryMethod::valueOf));
notificationRule.setConfiguration(fromJson(configuration, NotificationRuleConfig.class));
return notificationRule;
}

4
dao/src/main/java/org/thingsboard/server/dao/model/sql/NotificationTargetEntity.java

@ -55,7 +55,7 @@ public class NotificationTargetEntity extends BaseSqlEntity<NotificationTarget>
public NotificationTargetEntity(NotificationTarget notificationTarget) {
setId(notificationTarget.getUuidId());
setCreatedTime(notificationTarget.getCreatedTime());
setTenantId(getUuid(notificationTarget.getTenantId()));
setTenantId(getTenantUuid(notificationTarget.getTenantId()));
setName(notificationTarget.getName());
setConfiguration(toJson(notificationTarget.getConfiguration()));
}
@ -65,7 +65,7 @@ public class NotificationTargetEntity extends BaseSqlEntity<NotificationTarget>
NotificationTarget notificationTarget = new NotificationTarget();
notificationTarget.setId(new NotificationTargetId(id));
notificationTarget.setCreatedTime(createdTime);
notificationTarget.setTenantId(createId(tenantId, TenantId::fromUUID));
notificationTarget.setTenantId(getTenantId(tenantId));
notificationTarget.setName(name);
notificationTarget.setConfiguration(fromJson(configuration, NotificationTargetConfig.class));
return notificationTarget;

4
dao/src/main/java/org/thingsboard/server/dao/model/sql/NotificationTemplateEntity.java

@ -58,7 +58,7 @@ public class NotificationTemplateEntity extends BaseSqlEntity<NotificationTempla
public NotificationTemplateEntity(NotificationTemplate notificationTemplate) {
setId(notificationTemplate.getUuidId());
setCreatedTime(notificationTemplate.getCreatedTime());
setTenantId(getUuid(notificationTemplate.getTenantId()));
setTenantId(getTenantUuid(notificationTemplate.getTenantId()));
setName(notificationTemplate.getName());
setNotificationType(notificationTemplate.getNotificationType());
setConfiguration(toJson(notificationTemplate.getConfiguration()));
@ -69,7 +69,7 @@ public class NotificationTemplateEntity extends BaseSqlEntity<NotificationTempla
NotificationTemplate notificationTemplate = new NotificationTemplate();
notificationTemplate.setId(new NotificationTemplateId(id));
notificationTemplate.setCreatedTime(createdTime);
notificationTemplate.setTenantId(createId(tenantId, TenantId::fromUUID));
notificationTemplate.setTenantId(getTenantId(tenantId));
notificationTemplate.setName(name);
notificationTemplate.setNotificationType(notificationType);
notificationTemplate.setConfiguration(fromJson(configuration, NotificationTemplateConfig.class));

9
dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationRequestService.java

@ -43,7 +43,9 @@ public class DefaultNotificationRequestService implements NotificationRequestSer
@Override
public NotificationRequest saveNotificationRequest(TenantId tenantId, NotificationRequest notificationRequest) {
notificationRequestValidator.validate(notificationRequest, NotificationRequest::getTenantId);
return notificationRequestDao.save(tenantId, notificationRequest);
NotificationRequest savedNotificationRequest = notificationRequestDao.save(tenantId, notificationRequest);
savedNotificationRequest.copyContext(notificationRequest);
return savedNotificationRequest;
}
@Override
@ -56,6 +58,11 @@ public class DefaultNotificationRequestService implements NotificationRequestSer
return notificationRequestDao.findByTenantIdAndPageLink(tenantId, pageLink);
}
@Override
public List<NotificationRequestId> findNotificationRequestsIdsByStatusAndRuleId(TenantId tenantId, NotificationRequestStatus requestStatus, NotificationRuleId ruleId) {
return notificationRequestDao.findIdsByRuleId(tenantId, requestStatus, ruleId);
}
@Override
public List<NotificationRequest> findNotificationRequestsByRuleIdAndOriginatorEntityId(TenantId tenantId, NotificationRuleId ruleId, EntityId originatorEntityId) {
return notificationRequestDao.findByRuleIdAndOriginatorEntityId(tenantId, ruleId, originatorEntityId);

8
dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationRuleService.java

@ -39,8 +39,8 @@ public class DefaultNotificationRuleService implements NotificationRuleService {
}
@Override
public NotificationRule findNotificationRuleById(TenantId tenantId, NotificationRuleId notificationRuleId) {
return notificationRuleDao.findById(tenantId, notificationRuleId.getId());
public NotificationRule findNotificationRuleById(TenantId tenantId, NotificationRuleId id) {
return notificationRuleDao.findById(tenantId, id.getId());
}
@Override
@ -49,8 +49,8 @@ public class DefaultNotificationRuleService implements NotificationRuleService {
}
@Override
public void deleteNotificationRule(TenantId tenantId, NotificationRuleId notificationRuleId) {
notificationRuleDao.removeById(tenantId, notificationRuleId.getId());
public void deleteNotificationRuleById(TenantId tenantId, NotificationRuleId id) {
notificationRuleDao.removeById(tenantId, id.getId());
}
private static class NotificationRuleValidator extends DataValidator<NotificationRule> {

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

@ -23,6 +23,7 @@ import org.thingsboard.server.common.data.id.CustomerId;
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.NotificationRequestStatus;
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.NotificationTargetConfig;
@ -30,7 +31,6 @@ import org.thingsboard.server.common.data.notification.targets.SingleUserNotific
import org.thingsboard.server.common.data.notification.targets.UserListNotificationTargetConfig;
import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.dao.service.DataValidator;
import org.thingsboard.server.dao.user.UserService;
import java.util.List;
@ -43,12 +43,12 @@ import java.util.stream.Collectors;
public class DefaultNotificationTargetService implements NotificationTargetService {
private final NotificationTargetDao notificationTargetDao;
private final NotificationRequestDao notificationRequestDao;
private final NotificationRuleDao notificationRuleDao;
private final UserService userService;
private final NotificationTargetValidator validator = new NotificationTargetValidator();
@Override
public NotificationTarget saveNotificationTarget(TenantId tenantId, NotificationTarget notificationTarget) {
validator.validate(notificationTarget, NotificationTarget::getTenantId);
return notificationTargetDao.save(tenantId, notificationTarget);
}
@ -63,15 +63,15 @@ public class DefaultNotificationTargetService implements NotificationTargetServi
}
@Override
public PageData<User> findRecipientsForNotificationTarget(TenantId tenantId, NotificationTargetId notificationTargetId, PageLink pageLink) {
public PageData<User> findRecipientsForNotificationTarget(TenantId tenantId, CustomerId customerId, NotificationTargetId notificationTargetId, PageLink pageLink) {
NotificationTarget notificationTarget = findNotificationTargetById(tenantId, notificationTargetId);
Objects.requireNonNull(notificationTarget, "Notification target [" + notificationTargetId + "] not found");
NotificationTargetConfig configuration = notificationTarget.getConfiguration();
return findRecipientsForNotificationTargetConfig(tenantId, configuration, pageLink);
return findRecipientsForNotificationTargetConfig(tenantId, customerId, configuration, pageLink);
}
@Override
public PageData<User> findRecipientsForNotificationTargetConfig(TenantId tenantId, NotificationTargetConfig targetConfig, PageLink pageLink) {
public PageData<User> findRecipientsForNotificationTargetConfig(TenantId tenantId, CustomerId customerId, NotificationTargetConfig targetConfig, PageLink pageLink) {
switch (targetConfig.getType()) {
case SINGLE_USER: {
UserId userId = new UserId(((SingleUserNotificationTargetConfig) targetConfig).getUserId());
@ -88,8 +88,14 @@ public class DefaultNotificationTargetService implements NotificationTargetServi
if (tenantId.equals(TenantId.SYS_TENANT_ID)) {
throw new IllegalArgumentException("Customer users target is not supported for system administrator");
}
CustomerId customerId = new CustomerId(((CustomerUsersNotificationTargetConfig) targetConfig).getCustomerId());
return userService.findCustomerUsers(tenantId, customerId, pageLink);
CustomerUsersNotificationTargetConfig customerUsersConfig = (CustomerUsersNotificationTargetConfig) targetConfig;
if (!customerUsersConfig.isGetCustomerIdFromOriginatorEntity()) {
customerId = new CustomerId(customerUsersConfig.getCustomerId());
}
if (customerId != null && !customerId.isNullUid()) {
return userService.findCustomerUsers(tenantId, customerId, pageLink);
}
break;
}
case ALL_USERS: {
if (!tenantId.equals(TenantId.SYS_TENANT_ID)) {
@ -103,16 +109,14 @@ public class DefaultNotificationTargetService implements NotificationTargetServi
}
@Override
public void deleteNotificationTarget(TenantId tenantId, NotificationTargetId notificationTargetId) {
notificationTargetDao.removeById(tenantId, notificationTargetId.getId());
}
private static class NotificationTargetValidator extends DataValidator<NotificationTarget> {
@Override
protected void validateDataImpl(TenantId tenantId, NotificationTarget notificationTarget) {
super.validateDataImpl(tenantId, notificationTarget);
public void deleteNotificationTargetById(TenantId tenantId, NotificationTargetId id) {
if (notificationRequestDao.existsByStatusAndTargetId(tenantId, NotificationRequestStatus.SCHEDULED, id)) {
throw new IllegalArgumentException("Notification target is referenced by scheduled notification request");
}
if (notificationRuleDao.existsByTargetId(tenantId, id)) {
throw new IllegalArgumentException("Notification target is being used in notification rule");
}
notificationTargetDao.removeById(tenantId, id.getId());
}
}

19
dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationTemplateService.java

@ -19,13 +19,18 @@ import lombok.RequiredArgsConstructor;
import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.id.NotificationTemplateId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.notification.NotificationRequestStatus;
import org.thingsboard.server.common.data.notification.template.NotificationTemplate;
import org.thingsboard.server.dao.entity.AbstractEntityService;
import java.util.Map;
@Service
@RequiredArgsConstructor
public class DefaultNotificationTemplateService implements NotificationTemplateService {
public class DefaultNotificationTemplateService extends AbstractEntityService implements NotificationTemplateService {
private final NotificationTemplateDao notificationTemplateDao;
private final NotificationRequestDao notificationRequestDao;
@Override
public NotificationTemplate findNotificationTemplateById(TenantId tenantId, NotificationTemplateId id) {
@ -39,7 +44,17 @@ public class DefaultNotificationTemplateService implements NotificationTemplateS
@Override
public void deleteNotificationTemplateById(TenantId tenantId, NotificationTemplateId id) {
notificationTemplateDao.removeById(tenantId, id.getId());
if (notificationRequestDao.existsByStatusAndTemplateId(tenantId, NotificationRequestStatus.SCHEDULED, id)) {
throw new IllegalArgumentException("Notification template is referenced by scheduled notification request");
}
try {
notificationTemplateDao.removeById(tenantId, id.getId());
} catch (Exception e) {
checkConstraintViolation(e, Map.of(
"fk_notification_rule_template_id", "Notification template is referenced by notification rule"
));
throw e;
}
}
}

10
dao/src/main/java/org/thingsboard/server/dao/notification/NotificationRequestDao.java

@ -18,6 +18,8 @@ package org.thingsboard.server.dao.notification;
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.NotificationTargetId;
import org.thingsboard.server.common.data.id.NotificationTemplateId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.notification.NotificationRequest;
import org.thingsboard.server.common.data.notification.NotificationRequestStats;
@ -32,10 +34,18 @@ public interface NotificationRequestDao extends Dao<NotificationRequest> {
PageData<NotificationRequest> findByTenantIdAndPageLink(TenantId tenantId, PageLink pageLink);
List<NotificationRequestId> findIdsByRuleId(TenantId tenantId, NotificationRequestStatus requestStatus, NotificationRuleId ruleId);
List<NotificationRequest> findByRuleIdAndOriginatorEntityId(TenantId tenantId, NotificationRuleId ruleId, EntityId originatorEntityId);
PageData<NotificationRequest> findAllByStatus(NotificationRequestStatus status, PageLink pageLink);
void updateStatsById(TenantId tenantId, NotificationRequestId notificationRequestId, NotificationRequestStats stats);
boolean existsByStatusAndTargetId(TenantId tenantId, NotificationRequestStatus status, NotificationTargetId targetId);
boolean existsByStatusAndTemplateId(TenantId tenantId, NotificationRequestStatus status, NotificationTemplateId templateId);
int removeAllByCreatedTimeBefore(long ts);
}

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

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.dao.notification;
import org.thingsboard.server.common.data.id.NotificationTargetId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.notification.rule.NotificationRule;
import org.thingsboard.server.common.data.page.PageData;
@ -25,4 +26,6 @@ public interface NotificationRuleDao extends Dao<NotificationRule> {
PageData<NotificationRule> findByTenantIdAndPageLink(TenantId tenantId, PageLink pageLink);
boolean existsByTargetId(TenantId tenantId, NotificationTargetId targetId);
}

27
dao/src/main/java/org/thingsboard/server/dao/service/ConstraintValidator.java

@ -20,6 +20,11 @@ import lombok.extern.slf4j.Slf4j;
import org.hibernate.validator.HibernateValidator;
import org.hibernate.validator.HibernateValidatorConfiguration;
import org.hibernate.validator.cfg.ConstraintMapping;
import org.hibernate.validator.internal.cfg.context.DefaultConstraintMapping;
import org.hibernate.validator.internal.engine.ConfigurationImpl;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.validation.beanvalidation.LocalValidatorFactoryBean;
import org.thingsboard.server.common.data.validation.Length;
import org.thingsboard.server.common.data.validation.NoXss;
import org.thingsboard.server.dao.exception.DataValidationException;
@ -27,10 +32,12 @@ import org.thingsboard.server.dao.exception.DataValidationException;
import javax.validation.Path;
import javax.validation.Validation;
import javax.validation.Validator;
import javax.validation.constraints.AssertTrue;
import java.util.List;
import java.util.stream.Collectors;
@Slf4j
@Configuration
public class ConstraintValidator {
private static Validator fieldsValidator;
@ -69,12 +76,26 @@ public class ConstraintValidator {
private static void initializeValidators() {
HibernateValidatorConfiguration validatorConfiguration = Validation.byProvider(HibernateValidator.class).configure();
ConstraintMapping constraintMapping = validatorConfiguration.createConstraintMapping();
constraintMapping.constraintDefinition(NoXss.class).validatedBy(NoXssValidator.class);
constraintMapping.constraintDefinition(Length.class).validatedBy(StringLengthValidator.class);
ConstraintMapping constraintMapping = getCustomConstraintMapping();
validatorConfiguration.addMapping(constraintMapping);
fieldsValidator = validatorConfiguration.buildValidatorFactory().getValidator();
}
@Bean
public LocalValidatorFactoryBean validatorFactoryBean() {
LocalValidatorFactoryBean localValidatorFactoryBean = new LocalValidatorFactoryBean();
localValidatorFactoryBean.setConfigurationInitializer(configuration -> {
((ConfigurationImpl) configuration).addMapping(getCustomConstraintMapping());
});
return localValidatorFactoryBean;
}
private static ConstraintMapping getCustomConstraintMapping() {
ConstraintMapping constraintMapping = new DefaultConstraintMapping();
constraintMapping.constraintDefinition(NoXss.class).validatedBy(NoXssValidator.class);
constraintMapping.constraintDefinition(Length.class).validatedBy(StringLengthValidator.class);
return constraintMapping;
}
}

1
dao/src/main/java/org/thingsboard/server/dao/service/NoXssValidator.java

@ -65,4 +65,5 @@ public class NoXssValidator implements ConstraintValidator<NoXss, Object> {
return false;
}
}
}

24
dao/src/main/java/org/thingsboard/server/dao/sql/notification/JpaNotificationRequestDao.java

@ -23,6 +23,8 @@ import org.thingsboard.server.common.data.EntityType;
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.NotificationTargetId;
import org.thingsboard.server.common.data.id.NotificationTemplateId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.notification.NotificationRequest;
import org.thingsboard.server.common.data.notification.NotificationRequestStats;
@ -37,6 +39,7 @@ import org.thingsboard.server.dao.util.SqlDao;
import java.util.List;
import java.util.UUID;
import java.util.stream.Collectors;
@Component
@SqlDao
@ -50,6 +53,12 @@ public class JpaNotificationRequestDao extends JpaAbstractDao<NotificationReques
return DaoUtil.toPageData(notificationRequestRepository.findByTenantId(tenantId.getId(), DaoUtil.toPageable(pageLink)));
}
@Override
public List<NotificationRequestId> findIdsByRuleId(TenantId tenantId, NotificationRequestStatus requestStatus, NotificationRuleId ruleId) {
return notificationRequestRepository.findAllIdsByStatusAndRuleId(requestStatus, ruleId.getId()).stream()
.map(NotificationRequestId::new).collect(Collectors.toList());
}
@Override
public List<NotificationRequest> findByRuleIdAndOriginatorEntityId(TenantId tenantId, NotificationRuleId ruleId, EntityId originatorEntityId) {
return DaoUtil.convertDataList(notificationRequestRepository.findAllByRuleIdAndOriginatorEntityTypeAndOriginatorEntityId(ruleId.getId(), originatorEntityId.getEntityType(), originatorEntityId.getId()));
@ -65,6 +74,21 @@ public class JpaNotificationRequestDao extends JpaAbstractDao<NotificationReques
notificationRequestRepository.updateStatsById(notificationRequestId.getId(), JacksonUtil.valueToTree(stats));
}
@Override
public boolean existsByStatusAndTargetId(TenantId tenantId, NotificationRequestStatus status, NotificationTargetId targetId) {
return notificationRequestRepository.existsByStatusAndTargetsContaining(status, targetId.getId().toString());
}
@Override
public boolean existsByStatusAndTemplateId(TenantId tenantId, NotificationRequestStatus status, NotificationTemplateId templateId) {
return notificationRequestRepository.existsByStatusAndTemplateId(status, templateId.getId());
}
@Override
public int removeAllByCreatedTimeBefore(long ts) {
return notificationRequestRepository.deleteAllByCreatedTimeBefore(ts);
}
@Override
protected Class<NotificationRequestEntity> getEntityClass() {
return NotificationRequestEntity.class;

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

@ -20,6 +20,7 @@ import lombok.RequiredArgsConstructor;
import org.springframework.data.jpa.repository.JpaRepository;
import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.id.NotificationTargetId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.notification.rule.NotificationRule;
import org.thingsboard.server.common.data.page.PageData;
@ -45,6 +46,11 @@ public class JpaNotificationRuleDao extends JpaAbstractDao<NotificationRuleEntit
Strings.nullToEmpty(pageLink.getTextSearch()), DaoUtil.toPageable(pageLink)));
}
@Override
public boolean existsByTargetId(TenantId tenantId, NotificationTargetId targetId) {
return notificationRuleRepository.existsByConfigurationContaining(targetId.getId().toString());
}
@Override
protected Class<NotificationRuleEntity> getEntityClass() {
return NotificationRuleEntity.class;

10
dao/src/main/java/org/thingsboard/server/dao/sql/notification/NotificationRequestRepository.java

@ -36,6 +36,10 @@ public interface NotificationRequestRepository extends JpaRepository<Notificatio
Page<NotificationRequestEntity> findByTenantId(UUID tenantId, Pageable pageable);
@Query("SELECT r.id FROM NotificationRequestEntity r WHERE r.status = :status AND r.ruleId = :ruleId")
List<UUID> findAllIdsByStatusAndRuleId(@Param("status") NotificationRequestStatus status,
@Param("ruleId") UUID ruleId);
List<NotificationRequestEntity> findAllByRuleIdAndOriginatorEntityTypeAndOriginatorEntityId(UUID ruleId, EntityType originatorEntityType, UUID originatorEntityId);
Page<NotificationRequestEntity> findAllByStatus(NotificationRequestStatus status, Pageable pageable);
@ -45,4 +49,10 @@ public interface NotificationRequestRepository extends JpaRepository<Notificatio
@Query("UPDATE NotificationRequestEntity r SET r.stats = :stats WHERE r.id = :id")
void updateStatsById(@Param("id") UUID id, @Param("stats") JsonNode stats);
boolean existsByStatusAndTargetsContaining(NotificationRequestStatus status, String targetIdStr);
boolean existsByStatusAndTemplateId(NotificationRequestStatus status, UUID templateId);
int deleteAllByCreatedTimeBefore(long ts);
}

2
dao/src/main/java/org/thingsboard/server/dao/sql/notification/NotificationRuleRepository.java

@ -28,4 +28,6 @@ public interface NotificationRuleRepository extends JpaRepository<NotificationRu
Page<NotificationRuleEntity> findByTenantIdAndNameContainingIgnoreCase(UUID tenantId, String searchText, Pageable pageable);
boolean existsByConfigurationContaining(String string);
}

5
dao/src/main/java/org/thingsboard/server/dao/sqlts/insert/sql/SqlPartitioningRepository.java

@ -79,8 +79,9 @@ public class SqlPartitioningRepository {
}
}
public void dropPartitionsBefore(String table, long ts, long partitionDurationMs) {
public long dropPartitionsBefore(String table, long ts, long partitionDurationMs) {
List<Long> partitions = fetchPartitions(table);
long lastDroppedPartitionEndTime = -1;
for (Long partitionStartTime : partitions) {
long partitionEndTime = getPartitionEndTime(partitionStartTime, partitionDurationMs);
if (partitionEndTime < ts) {
@ -88,11 +89,13 @@ public class SqlPartitioningRepository {
boolean success = detachAndDropPartition(table, partitionStartTime);
if (success) {
log.info("[{}] Detached expired partition: {}", table, partitionStartTime);
lastDroppedPartitionEndTime = Math.max(partitionEndTime, lastDroppedPartitionEndTime);
}
} else {
log.debug("[{}] Skipping valid partition: {}", table, partitionStartTime);
}
}
return lastDroppedPartitionEndTime;
}
public void cleanupPartitionsCache(String table, long expTime, long partitionDurationMs) {

14
dao/src/main/resources/sql/schema-entities.sql

@ -783,7 +783,7 @@ CREATE TABLE IF NOT EXISTS user_auth_settings (
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
);
@ -791,7 +791,7 @@ CREATE TABLE IF NOT EXISTS notification_target (
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
@ -800,7 +800,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,
@ -810,16 +810,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)
);

2
rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java

@ -272,6 +272,8 @@ public interface TbContext {
ListeningExecutor getExternalCallExecutor();
ListeningExecutor getNotificationExecutor();
MailService getMailService(boolean isSystem);
SmsService getSmsService();

1
rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/slack/SlackService.java

@ -16,6 +16,7 @@
package org.thingsboard.rule.engine.api.slack;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.notification.template.SlackConversation;
import java.util.List;

6
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbNotificationNode.java

@ -51,15 +51,15 @@ public class TbNotificationNode implements TbNode {
public void onMsg(TbContext ctx, TbMsg msg) throws ExecutionException, InterruptedException, TbNodeException {
NotificationRequest notificationRequest = NotificationRequest.builder()
.tenantId(ctx.getTenantId())
.targetId(config.getTargetId())
.targets(config.getTargets())
.templateId(config.getTemplateId())
.deliveryMethods(config.getDeliveryMethods())
.originatorType(NotificationOriginatorType.RULE_NODE)
.originatorEntityId(ctx.getSelfId())
.originatorEntityId(ctx.getSelfId()) // todo: duplicate originator from msg originator, set originator's customerId
.build();
notificationRequest.setTemplateContext(msg.getMetaData().getData());
DonAsynchron.withCallback(ctx.getDbCallbackExecutor().executeAsync(() -> {
DonAsynchron.withCallback(ctx.getNotificationExecutor().executeAsync(() -> {
return ctx.getNotificationManager().processNotificationRequest(ctx.getTenantId(), notificationRequest);
}),
r -> {

7
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbNotificationNodeConfiguration.java

@ -25,13 +25,12 @@ import org.thingsboard.server.common.data.notification.NotificationRequestConfig
import javax.validation.constraints.NotEmpty;
import javax.validation.constraints.NotNull;
import java.util.List;
import java.util.UUID;
@Data
public class TbNotificationNodeConfiguration implements NodeConfiguration<TbNotificationNodeConfiguration> {
@NotNull
private NotificationTargetId targetId;
@NotEmpty
private List<NotificationTargetId> targets;
@NotNull
private NotificationTemplateId templateId;
@NotEmpty
@ -41,8 +40,6 @@ public class TbNotificationNodeConfiguration implements NodeConfiguration<TbNoti
@Override
public TbNotificationNodeConfiguration defaultConfiguration() {
TbNotificationNodeConfiguration config = new TbNotificationNodeConfiguration();
config.setTargetId(new NotificationTargetId(UUID.randomUUID()));
config.setTemplateId(new NotificationTemplateId(UUID.randomUUID()));
config.setDeliveryMethods(List.of(NotificationDeliveryMethod.WEBSOCKET));
config.setAdditionalConfig(new NotificationRequestConfig());
return config;

2
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbSlackNode.java

@ -23,7 +23,7 @@ import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNode;
import org.thingsboard.rule.engine.api.TbNodeConfiguration;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.rule.engine.api.slack.SlackConversation;
import org.thingsboard.server.common.data.notification.template.SlackConversation;
import org.thingsboard.rule.engine.api.util.TbNodeUtils;
import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg;

2
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbSlackNodeConfiguration.java

@ -17,7 +17,7 @@ package org.thingsboard.rule.engine.notification;
import lombok.Data;
import org.thingsboard.rule.engine.api.NodeConfiguration;
import org.thingsboard.rule.engine.api.slack.SlackConversation;
import org.thingsboard.server.common.data.notification.template.SlackConversation;
import javax.validation.constraints.NotEmpty;
import javax.validation.constraints.NotNull;

Loading…
Cancel
Save