diff --git a/application/src/main/data/upgrade/3.4.2/schema_update.sql b/application/src/main/data/upgrade/3.4.2/schema_update.sql index 7294916445..01113f2324 100644 --- a/application/src/main/data/upgrade/3.4.2/schema_update.sql +++ b/application/src/main/data/upgrade/3.4.2/schema_update.sql @@ -14,7 +14,6 @@ -- limitations under the License. -- - CREATE TABLE IF NOT EXISTS notification_target ( id UUID NOT NULL CONSTRAINT notification_target_pkey PRIMARY KEY, created_time BIGINT NOT NULL, @@ -22,12 +21,14 @@ CREATE TABLE IF NOT EXISTS notification_target ( name VARCHAR(255) NOT NULL, configuration varchar(1000) NOT NULL ); -CREATE INDEX IF NOT EXISTS idx_notification_target_tenant_id_and_created_time ON notification_target(tenant_id, created_time DESC); +CREATE INDEX IF NOT EXISTS idx_notification_target_tenant_id_created_time ON notification_target(tenant_id, created_time DESC); CREATE TABLE IF NOT EXISTS notification_rule ( id UUID NOT NULL CONSTRAINT notification_rule_pkey PRIMARY KEY ); +ALTER TABLE alarm ADD COLUMN IF NOT EXISTS notification_rule_id UUID; + CREATE TABLE IF NOT EXISTS notification_request ( id UUID NOT NULL CONSTRAINT notification_request_pkey PRIMARY KEY, created_time BIGINT NOT NULL, @@ -37,12 +38,14 @@ CREATE TABLE IF NOT EXISTS notification_request ( text_template VARCHAR NOT NULL, notification_info VARCHAR(1000), notification_severity VARCHAR(32), - additional_config VARCHAR(1000), - status VARCHAR(32), + 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), - alarm_id UUID + additional_config VARCHAR(1000), + status VARCHAR(32) ); -CREATE INDEX IF NOT EXISTS idx_notification_request_tenant_id_and_created_time ON notification_request(tenant_id, created_time DESC); +CREATE INDEX IF NOT EXISTS idx_notification_request_tenant_id_originator_type_created_time ON notification_request(tenant_id, originator_type, created_time DESC); CREATE TABLE IF NOT EXISTS notification ( id UUID NOT NULL, @@ -54,8 +57,9 @@ CREATE TABLE IF NOT EXISTS notification ( text VARCHAR NOT NULL, info VARCHAR(1000), severity VARCHAR(32), + originator_type VARCHAR(32) NOT NULL, status VARCHAR(32) ) PARTITION BY RANGE (created_time); CREATE INDEX IF NOT EXISTS idx_notification_id ON notification(id); +CREATE INDEX IF NOT EXISTS idx_notification_recipient_id_created_time ON notification(recipient_id, created_time DESC); CREATE INDEX IF NOT EXISTS idx_notification_notification_request_id ON notification(request_id); -CREATE INDEX IF NOT EXISTS idx_notification_recipient_id_and_created_time ON notification(recipient_id, created_time DESC); diff --git a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java index 1bfb515811..52cd33a885 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java +++ b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java @@ -30,7 +30,7 @@ import org.springframework.data.redis.core.RedisTemplate; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Component; import org.thingsboard.rule.engine.api.MailService; -import org.thingsboard.rule.engine.api.RuleEngineNotificationService; +import org.thingsboard.rule.engine.api.NotificationManager; import org.thingsboard.rule.engine.api.SmsService; import org.thingsboard.rule.engine.api.sms.SmsSenderFactory; import org.thingsboard.script.api.js.JsInvokeService; @@ -309,7 +309,7 @@ public class ActorSystemContext { @Autowired @Getter - private RuleEngineNotificationService notificationService; + private NotificationManager notificationManager; @Lazy @Autowired(required = false) diff --git a/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java b/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java index f623a9ad4c..bc43e69284 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java +++ b/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java @@ -25,10 +25,10 @@ import org.bouncycastle.util.Arrays; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.ListeningExecutor; import org.thingsboard.rule.engine.api.MailService; +import org.thingsboard.rule.engine.api.NotificationManager; import org.thingsboard.rule.engine.api.RuleEngineAlarmService; import org.thingsboard.rule.engine.api.RuleEngineAssetProfileCache; import org.thingsboard.rule.engine.api.RuleEngineDeviceProfileCache; -import org.thingsboard.rule.engine.api.RuleEngineNotificationService; import org.thingsboard.rule.engine.api.RuleEngineRpcService; import org.thingsboard.rule.engine.api.RuleEngineTelemetryService; import org.thingsboard.rule.engine.api.ScriptEngine; @@ -665,8 +665,8 @@ class DefaultTbContext implements TbContext { } @Override - public RuleEngineNotificationService getNotificationService() { - return mainCtx.getNotificationService(); + public NotificationManager getNotificationManager() { + return mainCtx.getNotificationManager(); } @Override diff --git a/application/src/main/java/org/thingsboard/server/controller/NotificationController.java b/application/src/main/java/org/thingsboard/server/controller/NotificationController.java index 5c621477a1..a392e6ec74 100644 --- a/application/src/main/java/org/thingsboard/server/controller/NotificationController.java +++ b/application/src/main/java/org/thingsboard/server/controller/NotificationController.java @@ -17,6 +17,7 @@ package org.thingsboard.server.controller; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; +import org.apache.commons.lang3.StringUtils; import org.springframework.security.access.prepost.PreAuthorize; import org.springframework.security.core.annotation.AuthenticationPrincipal; import org.springframework.web.bind.annotation.DeleteMapping; @@ -34,13 +35,15 @@ import org.thingsboard.server.common.data.exception.ThingsboardException; import org.thingsboard.server.common.data.id.NotificationId; import org.thingsboard.server.common.data.id.NotificationRequestId; import org.thingsboard.server.common.data.notification.Notification; +import org.thingsboard.server.common.data.notification.NotificationOriginatorType; import org.thingsboard.server.common.data.notification.NotificationRequest; +import org.thingsboard.server.common.data.notification.NotificationSeverity; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.dao.notification.NotificationRequestService; import org.thingsboard.server.dao.notification.NotificationService; import org.thingsboard.server.queue.util.TbCoreComponent; -import org.thingsboard.server.service.notification.NotificationSubscriptionService; +import org.thingsboard.rule.engine.api.NotificationManager; import org.thingsboard.server.service.security.model.SecurityUser; import org.thingsboard.server.service.security.permission.Operation; import org.thingsboard.server.service.security.permission.Resource; @@ -56,7 +59,7 @@ public class NotificationController extends BaseController { private final NotificationService notificationService; private final NotificationRequestService notificationRequestService; - private final NotificationSubscriptionService notificationSubscriptionService; + private final NotificationManager notificationManager; @GetMapping("/notifications") @PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN', 'CUSTOMER_USER')") @@ -76,7 +79,7 @@ public class NotificationController extends BaseController { public void markNotificationAsRead(@PathVariable UUID id, @AuthenticationPrincipal SecurityUser user) { NotificationId notificationId = new NotificationId(id); - notificationSubscriptionService.markNotificationAsRead(user.getTenantId(), user.getId(), notificationId); + notificationManager.markNotificationAsRead(user.getTenantId(), user.getId(), notificationId); } // delete notification? @@ -86,11 +89,27 @@ public class NotificationController extends BaseController { public NotificationRequest createNotificationRequest(@RequestBody NotificationRequest notificationRequest, @AuthenticationPrincipal SecurityUser user) throws ThingsboardException { accessControlService.checkPermission(user, Resource.NOTIFICATION_REQUEST, Operation.CREATE, null, notificationRequest); + // todo: check permission for notification target if (notificationRequest.getId() != null) { + // TODO: think about notification request update throw new IllegalArgumentException("Notification request cannot be changed. You can delete it and create a new one"); } + notificationRequest.setOriginatorType(NotificationOriginatorType.USER); + notificationRequest.setOriginatorEntityId(user.getId()); + if (StringUtils.isBlank(notificationRequest.getNotificationReason())) { + notificationRequest.setNotificationReason(NotificationRequest.GENERAL_NOTIFICATION_REASON); + } + if (notificationRequest.getNotificationSeverity() == null) { + notificationRequest.setNotificationSeverity(NotificationSeverity.NORMAL); + } + if (notificationRequest.getNotificationInfo() != null && notificationRequest.getNotificationInfo().getOriginatorType() != null) { + throw new IllegalArgumentException("Unsupported notification info type"); + } + notificationRequest.setRuleId(null); + notificationRequest.setStatus(null); + try { - NotificationRequest savedNotificationRequest = notificationSubscriptionService.processNotificationRequest(user.getTenantId(), notificationRequest); + NotificationRequest savedNotificationRequest = notificationManager.processNotificationRequest(user.getTenantId(), notificationRequest); logEntityAction(user, EntityType.NOTIFICATION_REQUEST, savedNotificationRequest, ActionType.ADDED); return savedNotificationRequest; } catch (Exception e) { @@ -126,7 +145,7 @@ public class NotificationController extends BaseController { NotificationRequest notificationRequest = notificationRequestService.findNotificationRequestById(user.getTenantId(), notificationRequestId); accessControlService.checkPermission(user, Resource.NOTIFICATION_REQUEST, Operation.DELETE, notificationRequestId, notificationRequest); try { - notificationSubscriptionService.deleteNotificationRequest(user.getTenantId(), notificationRequestId); + notificationManager.deleteNotificationRequest(user.getTenantId(), notificationRequestId); logEntityAction(user, EntityType.NOTIFICATION_REQUEST, notificationRequest, ActionType.DELETED); } catch (Exception e) { logEntityAction(user, EntityType.NOTIFICATION_REQUEST, notificationRequest, notificationRequest, ActionType.DELETED, e); diff --git a/application/src/main/java/org/thingsboard/server/controller/plugin/TbWebSocketHandler.java b/application/src/main/java/org/thingsboard/server/controller/plugin/TbWebSocketHandler.java index 7cb0a8258e..edd6843c2f 100644 --- a/application/src/main/java/org/thingsboard/server/controller/plugin/TbWebSocketHandler.java +++ b/application/src/main/java/org/thingsboard/server/controller/plugin/TbWebSocketHandler.java @@ -42,7 +42,7 @@ import org.thingsboard.server.service.security.model.UserPrincipal; import org.thingsboard.server.service.ws.SessionEvent; import org.thingsboard.server.service.ws.WebSocketMsgEndpoint; import org.thingsboard.server.service.ws.WebSocketSessionType; -import org.thingsboard.server.service.ws.telemetry.WebSocketService; +import org.thingsboard.server.service.ws.WebSocketService; import org.thingsboard.server.service.ws.WebSocketSessionRef; import javax.websocket.RemoteEndpoint; @@ -59,7 +59,7 @@ import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.LinkedBlockingQueue; -import static org.thingsboard.server.service.telemetry.DefaultWebSocketService.NUMBER_OF_PING_ATTEMPTS; +import static org.thingsboard.server.service.ws.DefaultWebSocketService.NUMBER_OF_PING_ATTEMPTS; @Service @TbCoreComponent diff --git a/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationSubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationManager.java similarity index 92% rename from application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationSubscriptionService.java rename to application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationManager.java index 8c786f33c9..4bb8afc61d 100644 --- a/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationSubscriptionService.java +++ b/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationManager.java @@ -18,7 +18,7 @@ package org.thingsboard.server.service.notification; import com.google.common.base.Strings; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; -import org.thingsboard.rule.engine.api.RuleEngineNotificationService; +import org.thingsboard.rule.engine.api.NotificationManager; import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.common.data.User; @@ -54,7 +54,7 @@ import java.util.UUID; @Service @Slf4j -public class DefaultNotificationSubscriptionService extends AbstractSubscriptionService implements NotificationSubscriptionService, RuleEngineNotificationService { +public class DefaultNotificationManager extends AbstractSubscriptionService implements NotificationManager { private final NotificationTargetService notificationTargetService; private final NotificationRequestService notificationRequestService; @@ -62,12 +62,12 @@ public class DefaultNotificationSubscriptionService extends AbstractSubscription private final DbCallbackExecutorService dbCallbackExecutorService; private final NotificationsTopicService notificationsTopicService; - public DefaultNotificationSubscriptionService(TbClusterService clusterService, PartitionService partitionService, - NotificationTargetService notificationTargetService, - NotificationRequestService notificationRequestService, - NotificationService notificationService, - DbCallbackExecutorService dbCallbackExecutorService, - NotificationsTopicService notificationsTopicService) { + public DefaultNotificationManager(TbClusterService clusterService, PartitionService partitionService, + NotificationTargetService notificationTargetService, + NotificationRequestService notificationRequestService, + NotificationService notificationService, + DbCallbackExecutorService dbCallbackExecutorService, + NotificationsTopicService notificationsTopicService) { super(clusterService, partitionService); this.notificationTargetService = notificationTargetService; this.notificationRequestService = notificationRequestService; @@ -81,7 +81,6 @@ public class DefaultNotificationSubscriptionService extends AbstractSubscription log.info("Processing notification request (tenant id: {}, notification target id: {})", tenantId, notificationRequest.getTargetId()); notificationRequest.setTenantId(tenantId); if (notificationRequest.getAdditionalConfig() != null) { - // TODO: think about notification request update NotificationRequestConfig config = notificationRequest.getAdditionalConfig(); if (config.getSendingDelayInMinutes() > 0 && notificationRequest.getId() == null) { notificationRequest.setStatus(NotificationRequestStatus.SCHEDULED); @@ -169,6 +168,7 @@ public class DefaultNotificationSubscriptionService extends AbstractSubscription .text(formatNotificationText(notificationRequest.getTextTemplate(), recipient)) .info(notificationRequest.getNotificationInfo()) .severity(notificationRequest.getNotificationSeverity()) + .originatorType(notificationRequest.getOriginatorType()) .status(NotificationStatus.SENT) .build(); return notificationService.saveNotification(recipient.getTenantId(), notification); diff --git a/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationRuleProcessingService.java b/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationRuleProcessingService.java index a36b5cc8c2..ac7e875e03 100644 --- a/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationRuleProcessingService.java +++ b/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationRuleProcessingService.java @@ -19,11 +19,14 @@ import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import lombok.RequiredArgsConstructor; import org.springframework.stereotype.Service; +import org.thingsboard.rule.engine.api.NotificationManager; import org.thingsboard.server.common.data.alarm.Alarm; import org.thingsboard.server.common.data.id.NotificationRuleId; import org.thingsboard.server.common.data.id.NotificationTargetId; import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.notification.AlarmOriginatedNotificationInfo; 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; @@ -44,7 +47,7 @@ public class DefaultNotificationRuleProcessingService implements NotificationRul private final NotificationRuleService notificationRuleService; private final NotificationRequestService notificationRequestService; - private final NotificationSubscriptionService notificationSubscriptionService; + private final NotificationManager notificationManager; private final DbCallbackExecutorService dbCallbackExecutorService; @Override @@ -65,7 +68,7 @@ public class DefaultNotificationRuleProcessingService implements NotificationRul } private void onAlarmUpdate(TenantId tenantId, NotificationRuleId notificationRuleId, Alarm alarm) { - List notificationRequests = notificationRequestService.findNotificationRequestsByRuleIdAndAlarmId(tenantId, notificationRuleId, alarm.getId()); + List notificationRequests = notificationRequestService.findNotificationRequestsByRuleIdAndOriginatorEntityId(tenantId, notificationRuleId, alarm.getId()); NotificationRule notificationRule = notificationRuleService.findNotificationRuleById(tenantId, notificationRuleId); if (notificationRequests.isEmpty()) { // in case it is first notification for alarm, or it was previously acked and now we need to send notifications again @@ -82,7 +85,7 @@ public class DefaultNotificationRuleProcessingService implements NotificationRul // using regular service due to no need to send an update to subscription manager notificationRequestService.deleteNotificationRequestById(tenantId, notificationRequest.getId()); } else { - notificationSubscriptionService.deleteNotificationRequest(tenantId, notificationRequest.getId()); + notificationManager.deleteNotificationRequest(tenantId, notificationRequest.getId()); // todo: or should we mark already sent notifications as read? } } @@ -92,7 +95,7 @@ public class DefaultNotificationRuleProcessingService implements NotificationRul NotificationInfo previousNotificationInfo = notificationRequest.getNotificationInfo(); if (!previousNotificationInfo.equals(newNotificationInfo)) { notificationRequest.setNotificationInfo(newNotificationInfo); - notificationSubscriptionService.updateNotificationRequest(tenantId, notificationRequest); + notificationManager.updateNotificationRequest(tenantId, notificationRequest); } // fixme: no need to send an update event for scheduled requests, only for sent } @@ -118,15 +121,16 @@ public class DefaultNotificationRuleProcessingService implements NotificationRul .textTemplate(notificationRule.getNotificationTextTemplate()) // todo: format with alarm vars .notificationInfo(notificationInfo) .notificationSeverity(NotificationSeverity.NORMAL) // todo: from alarm severity - .additionalConfig(config) + .originatorType(NotificationOriginatorType.ALARM) + .originatorEntityId(alarm.getId()) .ruleId(notificationRule.getId()) - .alarmId(alarm.getId()) + .additionalConfig(config) .build(); - notificationSubscriptionService.processNotificationRequest(tenantId, notificationRequest); + notificationManager.processNotificationRequest(tenantId, notificationRequest); } private NotificationInfo constructNotificationInfo(Alarm alarm, NotificationRule notificationRule) { - return NotificationInfo.builder() + return AlarmOriginatedNotificationInfo.builder() .alarmId(alarm.getId()) .alarmType(alarm.getType()) .alarmOriginator(alarm.getOriginator()) diff --git a/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationSchedulerService.java b/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationSchedulerService.java index 903444529b..36faeeee87 100644 --- a/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationSchedulerService.java +++ b/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationSchedulerService.java @@ -20,6 +20,7 @@ import com.google.common.util.concurrent.ListenableScheduledFuture; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; +import org.thingsboard.rule.engine.api.NotificationManager; import org.thingsboard.server.common.data.id.NotificationRequestId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.notification.NotificationRequest; @@ -48,7 +49,7 @@ import java.util.concurrent.TimeUnit; @SuppressWarnings("UnstableApiUsage") public class DefaultNotificationSchedulerService extends AbstractPartitionBasedService implements NotificationSchedulerService { - private final NotificationSubscriptionService notificationSubscriptionService; + private final NotificationManager notificationManager; private final NotificationRequestService notificationRequestService; private final Map> scheduledNotificationRequests = new ConcurrentHashMap<>(); @@ -93,7 +94,7 @@ public class DefaultNotificationSchedulerService extends AbstractPartitionBasedS NotificationRequest notificationRequest = notificationRequestService.findNotificationRequestById(tenantId, request.getId()); if (notificationRequest == null) return; - notificationSubscriptionService.processNotificationRequest(tenantId, notificationRequest); + notificationManager.processNotificationRequest(tenantId, notificationRequest); scheduledNotificationRequests.remove(notificationRequest.getId()); }, delayMs, TimeUnit.MILLISECONDS); scheduledNotificationRequests.put(request.getId(), scheduledTask); diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java index eed12acf3e..4661730f9c 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java @@ -48,7 +48,7 @@ import org.thingsboard.server.dao.timeseries.TimeseriesService; import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.executors.DbCallbackExecutorService; -import org.thingsboard.server.service.ws.telemetry.WebSocketService; +import org.thingsboard.server.service.ws.WebSocketService; import org.thingsboard.server.service.ws.WebSocketSessionRef; import org.thingsboard.server.service.ws.telemetry.cmd.v2.AggHistoryCmd; import org.thingsboard.server.service.ws.telemetry.cmd.v2.AggKey; diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractDataSubCtx.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractDataSubCtx.java index 7ff8bf931e..11fa6b4c35 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractDataSubCtx.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractDataSubCtx.java @@ -29,7 +29,7 @@ import org.thingsboard.server.common.data.query.TsValue; import org.thingsboard.server.dao.attributes.AttributesService; import org.thingsboard.server.dao.entity.EntityService; import org.thingsboard.server.service.ws.WebSocketSessionRef; -import org.thingsboard.server.service.ws.telemetry.WebSocketService; +import org.thingsboard.server.service.ws.WebSocketService; import org.thingsboard.server.service.ws.telemetry.sub.TelemetrySubscriptionUpdate; import java.util.ArrayList; diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractSubCtx.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractSubCtx.java index d7fde23b31..8dd90da5dc 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractSubCtx.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractSubCtx.java @@ -39,7 +39,7 @@ import org.thingsboard.server.common.data.query.TsValue; import org.thingsboard.server.dao.attributes.AttributesService; import org.thingsboard.server.dao.entity.EntityService; import org.thingsboard.server.service.ws.WebSocketSessionRef; -import org.thingsboard.server.service.ws.telemetry.WebSocketService; +import org.thingsboard.server.service.ws.WebSocketService; import org.thingsboard.server.service.ws.telemetry.cmd.v2.CmdUpdate; import org.thingsboard.server.service.ws.telemetry.sub.TelemetrySubscriptionUpdate; diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbAlarmDataSubCtx.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbAlarmDataSubCtx.java index c66cdc4828..a952a01ee9 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/TbAlarmDataSubCtx.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbAlarmDataSubCtx.java @@ -38,7 +38,7 @@ import org.thingsboard.server.dao.alarm.AlarmService; import org.thingsboard.server.dao.attributes.AttributesService; import org.thingsboard.server.dao.entity.EntityService; import org.thingsboard.server.dao.model.ModelConstants; -import org.thingsboard.server.service.ws.telemetry.WebSocketService; +import org.thingsboard.server.service.ws.WebSocketService; import org.thingsboard.server.service.ws.WebSocketSessionRef; import org.thingsboard.server.service.ws.telemetry.cmd.v2.AlarmDataUpdate; import org.thingsboard.server.service.ws.telemetry.sub.AlarmSubscriptionUpdate; diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityCountSubCtx.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityCountSubCtx.java index f0cc2260b6..573016c20c 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityCountSubCtx.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityCountSubCtx.java @@ -19,7 +19,7 @@ import lombok.extern.slf4j.Slf4j; import org.thingsboard.server.common.data.query.EntityCountQuery; import org.thingsboard.server.dao.attributes.AttributesService; import org.thingsboard.server.dao.entity.EntityService; -import org.thingsboard.server.service.ws.telemetry.WebSocketService; +import org.thingsboard.server.service.ws.WebSocketService; import org.thingsboard.server.service.ws.WebSocketSessionRef; import org.thingsboard.server.service.ws.telemetry.cmd.v2.EntityCountUpdate; diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java index 975ecea1c2..f5c701916d 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java @@ -27,7 +27,7 @@ import org.thingsboard.server.common.data.query.EntityKeyType; import org.thingsboard.server.common.data.query.TsValue; import org.thingsboard.server.dao.attributes.AttributesService; import org.thingsboard.server.dao.entity.EntityService; -import org.thingsboard.server.service.ws.telemetry.WebSocketService; +import org.thingsboard.server.service.ws.WebSocketService; import org.thingsboard.server.service.ws.WebSocketSessionRef; import org.thingsboard.server.service.ws.telemetry.cmd.v2.EntityDataCmd; import org.thingsboard.server.service.ws.telemetry.cmd.v2.EntityDataUpdate; diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultWebSocketService.java b/application/src/main/java/org/thingsboard/server/service/ws/DefaultWebSocketService.java similarity index 95% rename from application/src/main/java/org/thingsboard/server/service/telemetry/DefaultWebSocketService.java rename to application/src/main/java/org/thingsboard/server/service/ws/DefaultWebSocketService.java index e03e0fc3fb..2598429672 100644 --- a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultWebSocketService.java +++ b/application/src/main/java/org/thingsboard/server/service/ws/DefaultWebSocketService.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.service.telemetry; +package org.thingsboard.server.service.ws; import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.databind.ObjectMapper; @@ -62,14 +62,9 @@ import org.thingsboard.server.service.subscription.TbAttributeSubscriptionScope; import org.thingsboard.server.service.subscription.TbEntityDataSubscriptionService; import org.thingsboard.server.service.subscription.TbLocalSubscriptionService; import org.thingsboard.server.service.subscription.TbTimeseriesSubscription; -import org.thingsboard.server.service.ws.SessionEvent; -import org.thingsboard.server.service.ws.WebSocketMsgEndpoint; -import org.thingsboard.server.service.ws.WebSocketSessionRef; -import org.thingsboard.server.service.ws.WsCmd; -import org.thingsboard.server.service.ws.WsSessionMetaData; import org.thingsboard.server.service.ws.notification.NotificationCommandsHandler; import org.thingsboard.server.service.ws.notification.cmd.NotificationCmdsWrapper; -import org.thingsboard.server.service.ws.telemetry.WebSocketService; +import org.thingsboard.server.service.ws.notification.cmd.WsCmd; import org.thingsboard.server.service.ws.telemetry.cmd.TelemetryPluginCmdsWrapper; import org.thingsboard.server.service.ws.telemetry.cmd.v1.AttributesSubscriptionCmd; import org.thingsboard.server.service.ws.telemetry.cmd.v1.GetHistoryCmd; @@ -163,22 +158,22 @@ public class DefaultWebSocketService implements WebSocketService { pingExecutor.scheduleWithFixedDelay(this::sendPing, pingTimeout / NUMBER_OF_PING_ATTEMPTS, pingTimeout / NUMBER_OF_PING_ATTEMPTS, TimeUnit.MILLISECONDS); telemetryCmdsHandlers = List.of( - WsCmdListHandler.of(TelemetryPluginCmdsWrapper::getAttrSubCmds, this::handleWsAttributesSubscriptionCmd), - WsCmdListHandler.of(TelemetryPluginCmdsWrapper::getTsSubCmds, this::handleWsTimeseriesSubscriptionCmd), - WsCmdListHandler.of(TelemetryPluginCmdsWrapper::getHistoryCmds, this::handleWsHistoryCmd), - WsCmdListHandler.of(TelemetryPluginCmdsWrapper::getEntityDataCmds, this::handleWsEntityDataCmd), - WsCmdListHandler.of(TelemetryPluginCmdsWrapper::getAlarmDataCmds, this::handleWsAlarmDataCmd), - WsCmdListHandler.of(TelemetryPluginCmdsWrapper::getEntityCountCmds, this::handleWsEntityCountCmd), - WsCmdListHandler.of(TelemetryPluginCmdsWrapper::getEntityDataUnsubscribeCmds, this::handleWsDataUnsubscribeCmd), - WsCmdListHandler.of(TelemetryPluginCmdsWrapper::getAlarmDataUnsubscribeCmds, this::handleWsDataUnsubscribeCmd), - WsCmdListHandler.of(TelemetryPluginCmdsWrapper::getAlarmDataUnsubscribeCmds, this::handleWsDataUnsubscribeCmd), - WsCmdListHandler.of(TelemetryPluginCmdsWrapper::getEntityCountUnsubscribeCmds, this::handleWsDataUnsubscribeCmd) + newCmdsHandler(TelemetryPluginCmdsWrapper::getAttrSubCmds, this::handleWsAttributesSubscriptionCmd), + newCmdsHandler(TelemetryPluginCmdsWrapper::getTsSubCmds, this::handleWsTimeseriesSubscriptionCmd), + newCmdsHandler(TelemetryPluginCmdsWrapper::getHistoryCmds, this::handleWsHistoryCmd), + newCmdsHandler(TelemetryPluginCmdsWrapper::getEntityDataCmds, this::handleWsEntityDataCmd), + newCmdsHandler(TelemetryPluginCmdsWrapper::getAlarmDataCmds, this::handleWsAlarmDataCmd), + newCmdsHandler(TelemetryPluginCmdsWrapper::getEntityCountCmds, this::handleWsEntityCountCmd), + newCmdsHandler(TelemetryPluginCmdsWrapper::getEntityDataUnsubscribeCmds, this::handleWsDataUnsubscribeCmd), + newCmdsHandler(TelemetryPluginCmdsWrapper::getAlarmDataUnsubscribeCmds, this::handleWsDataUnsubscribeCmd), + newCmdsHandler(TelemetryPluginCmdsWrapper::getAlarmDataUnsubscribeCmds, this::handleWsDataUnsubscribeCmd), + newCmdsHandler(TelemetryPluginCmdsWrapper::getEntityCountUnsubscribeCmds, this::handleWsDataUnsubscribeCmd) ); notificationCmdsHandlers = List.of( - WsCmdHandler.of(NotificationCmdsWrapper::getUnreadSubCmd, notificationCmdsHandler::handleUnreadNotificationsSubCmd), - WsCmdHandler.of(NotificationCmdsWrapper::getUnreadCountSubCmd, notificationCmdsHandler::handleUnreadNotificationsCountSubCmd), - WsCmdHandler.of(NotificationCmdsWrapper::getMarkAsReadCmd, notificationCmdsHandler::handleMarkAsReadCmd), - WsCmdHandler.of(NotificationCmdsWrapper::getUnsubCmd, notificationCmdsHandler::handleUnsubCmd) + newCmdHandler(NotificationCmdsWrapper::getUnreadSubCmd, notificationCmdsHandler::handleUnreadNotificationsSubCmd), + newCmdHandler(NotificationCmdsWrapper::getUnreadCountSubCmd, notificationCmdsHandler::handleUnreadNotificationsCountSubCmd), + newCmdHandler(NotificationCmdsWrapper::getMarkAsReadCmd, notificationCmdsHandler::handleMarkAsReadCmd), + newCmdHandler(NotificationCmdsWrapper::getUnsubCmd, notificationCmdsHandler::handleUnsubCmd) ); } @@ -955,7 +950,17 @@ public class DefaultWebSocketService implements WebSocketService { return limit == 0 ? DEFAULT_LIMIT : limit; } - @RequiredArgsConstructor(staticName = "of") + public static WsCmdHandler newCmdHandler(java.util.function.Function cmdExtractor, + BiConsumer handler) { + return new WsCmdHandler<>(cmdExtractor, handler); + } + + public static WsCmdListHandler newCmdsHandler(java.util.function.Function> cmdsExtractor, + BiConsumer handler) { + return new WsCmdListHandler<>(cmdsExtractor, handler); + } + + @RequiredArgsConstructor public static class WsCmdHandler { private final java.util.function.Function cmdExtractor; private final BiConsumer handler; @@ -970,13 +975,13 @@ public class DefaultWebSocketService implements WebSocketService { } } - @RequiredArgsConstructor(staticName = "of") + @RequiredArgsConstructor public static class WsCmdListHandler { - private final java.util.function.Function> cmdExtractor; + private final java.util.function.Function> cmdsExtractor; private final BiConsumer handler; public List extractCmds(W cmdsWrapper) { - return cmdExtractor.apply(cmdsWrapper); + return cmdsExtractor.apply(cmdsWrapper); } @SuppressWarnings("unchecked") diff --git a/application/src/main/java/org/thingsboard/server/service/ws/telemetry/WebSocketService.java b/application/src/main/java/org/thingsboard/server/service/ws/WebSocketService.java similarity index 96% rename from application/src/main/java/org/thingsboard/server/service/ws/telemetry/WebSocketService.java rename to application/src/main/java/org/thingsboard/server/service/ws/WebSocketService.java index d0d8dab59e..ddd08590b5 100644 --- a/application/src/main/java/org/thingsboard/server/service/ws/telemetry/WebSocketService.java +++ b/application/src/main/java/org/thingsboard/server/service/ws/WebSocketService.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.service.ws.telemetry; +package org.thingsboard.server.service.ws; import org.springframework.web.socket.CloseStatus; import org.thingsboard.server.service.ws.telemetry.cmd.v2.CmdUpdate; diff --git a/application/src/main/java/org/thingsboard/server/service/ws/notification/DefaultNotificationCommandsHandler.java b/application/src/main/java/org/thingsboard/server/service/ws/notification/DefaultNotificationCommandsHandler.java index 6af03caeb5..5e31924e45 100644 --- a/application/src/main/java/org/thingsboard/server/service/ws/notification/DefaultNotificationCommandsHandler.java +++ b/application/src/main/java/org/thingsboard/server/service/ws/notification/DefaultNotificationCommandsHandler.java @@ -31,7 +31,7 @@ import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.dao.notification.NotificationService; import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; import org.thingsboard.server.queue.util.TbCoreComponent; -import org.thingsboard.server.service.notification.NotificationSubscriptionService; +import org.thingsboard.rule.engine.api.NotificationManager; import org.thingsboard.server.service.security.model.SecurityUser; import org.thingsboard.server.service.ws.notification.cmd.NotificationsCountSubCmd; import org.thingsboard.server.service.ws.notification.sub.NotificationRequestUpdate; @@ -39,11 +39,11 @@ import org.thingsboard.server.service.ws.notification.sub.NotificationUpdate; import org.thingsboard.server.service.ws.notification.sub.NotificationsSubscription; import org.thingsboard.server.service.subscription.TbLocalSubscriptionService; import org.thingsboard.server.service.ws.WebSocketSessionRef; -import org.thingsboard.server.service.ws.notification.cmd.MarkNotificationAsReadCmd; +import org.thingsboard.server.service.ws.notification.cmd.MarkNotificationsAsReadCmd; import org.thingsboard.server.service.ws.notification.cmd.NotificationsSubCmd; import org.thingsboard.server.service.ws.notification.sub.NotificationsSubscriptionUpdate; import org.thingsboard.server.service.ws.notification.sub.NotificationsCountSubscription; -import org.thingsboard.server.service.ws.telemetry.WebSocketService; +import org.thingsboard.server.service.ws.WebSocketService; import org.thingsboard.server.service.ws.telemetry.cmd.v2.CmdUpdate; import org.thingsboard.server.service.ws.telemetry.cmd.v2.UnsubscribeCmd; @@ -59,7 +59,7 @@ public class DefaultNotificationCommandsHandler implements NotificationCommandsH private final NotificationService notificationService; private final TbLocalSubscriptionService localSubscriptionService; - private final NotificationSubscriptionService notificationSubscriptionService; + private final NotificationManager notificationManager; private final TbServiceInfoProvider serviceInfoProvider; @Autowired @Lazy private WebSocketService wsService; @@ -67,13 +67,13 @@ public class DefaultNotificationCommandsHandler implements NotificationCommandsH @Override public void handleUnreadNotificationsSubCmd(WebSocketSessionRef sessionRef, NotificationsSubCmd cmd) { log.debug("[{}] Handling unread notifications subscription cmd (cmdId: {})", sessionRef.getSessionId(), cmd.getCmdId()); - SecurityUser user = sessionRef.getSecurityCtx(); + SecurityUser securityCtx = sessionRef.getSecurityCtx(); NotificationsSubscription subscription = NotificationsSubscription.builder() .serviceId(serviceInfoProvider.getServiceId()) .sessionId(sessionRef.getSessionId()) .subscriptionId(cmd.getCmdId()) - .tenantId(user.getTenantId()) - .entityId(user.getId()) + .tenantId(securityCtx.getTenantId()) + .entityId(securityCtx.getId()) .updateProcessor(this::handleNotificationsSubscriptionUpdate) .limit(cmd.getLimit()) .build(); @@ -86,13 +86,13 @@ public class DefaultNotificationCommandsHandler implements NotificationCommandsH @Override public void handleUnreadNotificationsCountSubCmd(WebSocketSessionRef sessionRef, NotificationsCountSubCmd cmd) { log.debug("[{}] Handling unread notifications count subscription cmd (cmdId: {})", sessionRef.getSessionId(), cmd.getCmdId()); - SecurityUser user = sessionRef.getSecurityCtx(); + SecurityUser securityCtx = sessionRef.getSecurityCtx(); NotificationsCountSubscription subscription = NotificationsCountSubscription.builder() .serviceId(serviceInfoProvider.getServiceId()) .sessionId(sessionRef.getSessionId()) .subscriptionId(cmd.getCmdId()) - .tenantId(user.getTenantId()) - .entityId(user.getId()) + .tenantId(securityCtx.getTenantId()) + .entityId(securityCtx.getId()) .updateProcessor(this::handleNotificationsCountSubscriptionUpdate) .build(); localSubscriptionService.addSubscription(subscription); @@ -203,9 +203,14 @@ public class DefaultNotificationCommandsHandler implements NotificationCommandsH @Override - public void handleMarkAsReadCmd(WebSocketSessionRef sessionRef, MarkNotificationAsReadCmd cmd) { - NotificationId notificationId = new NotificationId(cmd.getNotificationId()); - notificationSubscriptionService.markNotificationAsRead(sessionRef.getSecurityCtx().getTenantId(), sessionRef.getSecurityCtx().getId(), notificationId); + public void handleMarkAsReadCmd(WebSocketSessionRef sessionRef, MarkNotificationsAsReadCmd cmd) { + SecurityUser securityCtx = sessionRef.getSecurityCtx(); + cmd.getNotifications().stream() + .map(NotificationId::new) + .forEach(notificationId -> { + notificationManager.markNotificationAsRead(securityCtx.getTenantId(), securityCtx.getId(), notificationId); + // fixme: should send bulk update event, not a separate event for each notification + }); } @Override diff --git a/application/src/main/java/org/thingsboard/server/service/ws/notification/NotificationCommandsHandler.java b/application/src/main/java/org/thingsboard/server/service/ws/notification/NotificationCommandsHandler.java index 4cf1fb3809..0f57730601 100644 --- a/application/src/main/java/org/thingsboard/server/service/ws/notification/NotificationCommandsHandler.java +++ b/application/src/main/java/org/thingsboard/server/service/ws/notification/NotificationCommandsHandler.java @@ -16,7 +16,7 @@ package org.thingsboard.server.service.ws.notification; import org.thingsboard.server.service.ws.WebSocketSessionRef; -import org.thingsboard.server.service.ws.notification.cmd.MarkNotificationAsReadCmd; +import org.thingsboard.server.service.ws.notification.cmd.MarkNotificationsAsReadCmd; import org.thingsboard.server.service.ws.notification.cmd.NotificationsSubCmd; import org.thingsboard.server.service.ws.notification.cmd.NotificationsCountSubCmd; import org.thingsboard.server.service.ws.telemetry.cmd.v2.UnsubscribeCmd; @@ -27,7 +27,7 @@ public interface NotificationCommandsHandler { void handleUnreadNotificationsCountSubCmd(WebSocketSessionRef sessionRef, NotificationsCountSubCmd cmd); - void handleMarkAsReadCmd(WebSocketSessionRef sessionRef, MarkNotificationAsReadCmd cmd); + void handleMarkAsReadCmd(WebSocketSessionRef sessionRef, MarkNotificationsAsReadCmd cmd); void handleUnsubCmd(WebSocketSessionRef sessionRef, UnsubscribeCmd cmd); diff --git a/application/src/main/java/org/thingsboard/server/service/ws/notification/cmd/MarkNotificationAsReadCmd.java b/application/src/main/java/org/thingsboard/server/service/ws/notification/cmd/MarkNotificationsAsReadCmd.java similarity index 86% rename from application/src/main/java/org/thingsboard/server/service/ws/notification/cmd/MarkNotificationAsReadCmd.java rename to application/src/main/java/org/thingsboard/server/service/ws/notification/cmd/MarkNotificationsAsReadCmd.java index cfc6e8e38b..6ed49fc4b6 100644 --- a/application/src/main/java/org/thingsboard/server/service/ws/notification/cmd/MarkNotificationAsReadCmd.java +++ b/application/src/main/java/org/thingsboard/server/service/ws/notification/cmd/MarkNotificationsAsReadCmd.java @@ -18,14 +18,14 @@ package org.thingsboard.server.service.ws.notification.cmd; import lombok.AllArgsConstructor; import lombok.Data; import lombok.NoArgsConstructor; -import org.thingsboard.server.service.ws.WsCmd; +import java.util.List; import java.util.UUID; @Data @NoArgsConstructor @AllArgsConstructor -public class MarkNotificationAsReadCmd implements WsCmd { +public class MarkNotificationsAsReadCmd implements WsCmd { private int cmdId; - private UUID notificationId; + private List notifications; } diff --git a/application/src/main/java/org/thingsboard/server/service/ws/notification/cmd/NotificationCmdsWrapper.java b/application/src/main/java/org/thingsboard/server/service/ws/notification/cmd/NotificationCmdsWrapper.java index 011e8cecca..1c56d49a6f 100644 --- a/application/src/main/java/org/thingsboard/server/service/ws/notification/cmd/NotificationCmdsWrapper.java +++ b/application/src/main/java/org/thingsboard/server/service/ws/notification/cmd/NotificationCmdsWrapper.java @@ -23,7 +23,8 @@ public class NotificationCmdsWrapper { private NotificationsCountSubCmd unreadCountSubCmd; private NotificationsSubCmd unreadSubCmd; - private MarkNotificationAsReadCmd markAsReadCmd; + + private MarkNotificationsAsReadCmd markAsReadCmd; private NotificationsUnsubCmd unsubCmd; diff --git a/application/src/main/java/org/thingsboard/server/service/ws/notification/cmd/NotificationsCountSubCmd.java b/application/src/main/java/org/thingsboard/server/service/ws/notification/cmd/NotificationsCountSubCmd.java index 2722d6ae7d..985e0c9aa0 100644 --- a/application/src/main/java/org/thingsboard/server/service/ws/notification/cmd/NotificationsCountSubCmd.java +++ b/application/src/main/java/org/thingsboard/server/service/ws/notification/cmd/NotificationsCountSubCmd.java @@ -17,9 +17,7 @@ package org.thingsboard.server.service.ws.notification.cmd; import lombok.AllArgsConstructor; import lombok.Data; -import lombok.EqualsAndHashCode; import lombok.NoArgsConstructor; -import org.thingsboard.server.service.ws.WsCmd; @Data @NoArgsConstructor diff --git a/application/src/main/java/org/thingsboard/server/service/ws/notification/cmd/NotificationsSubCmd.java b/application/src/main/java/org/thingsboard/server/service/ws/notification/cmd/NotificationsSubCmd.java index 087b250314..c772635176 100644 --- a/application/src/main/java/org/thingsboard/server/service/ws/notification/cmd/NotificationsSubCmd.java +++ b/application/src/main/java/org/thingsboard/server/service/ws/notification/cmd/NotificationsSubCmd.java @@ -17,9 +17,7 @@ package org.thingsboard.server.service.ws.notification.cmd; import lombok.AllArgsConstructor; import lombok.Data; -import lombok.EqualsAndHashCode; import lombok.NoArgsConstructor; -import org.thingsboard.server.service.ws.WsCmd; @Data @NoArgsConstructor diff --git a/application/src/main/java/org/thingsboard/server/service/ws/notification/cmd/NotificationsUnsubCmd.java b/application/src/main/java/org/thingsboard/server/service/ws/notification/cmd/NotificationsUnsubCmd.java index b81e76e81d..d8e825e21d 100644 --- a/application/src/main/java/org/thingsboard/server/service/ws/notification/cmd/NotificationsUnsubCmd.java +++ b/application/src/main/java/org/thingsboard/server/service/ws/notification/cmd/NotificationsUnsubCmd.java @@ -18,7 +18,6 @@ package org.thingsboard.server.service.ws.notification.cmd; import lombok.AllArgsConstructor; import lombok.Data; import lombok.NoArgsConstructor; -import org.thingsboard.server.service.ws.WsCmd; import org.thingsboard.server.service.ws.telemetry.cmd.v2.UnsubscribeCmd; @Data diff --git a/application/src/main/java/org/thingsboard/server/service/ws/WsCmd.java b/application/src/main/java/org/thingsboard/server/service/ws/notification/cmd/WsCmd.java similarity index 91% rename from application/src/main/java/org/thingsboard/server/service/ws/WsCmd.java rename to application/src/main/java/org/thingsboard/server/service/ws/notification/cmd/WsCmd.java index 096c5930a5..daa59abcba 100644 --- a/application/src/main/java/org/thingsboard/server/service/ws/WsCmd.java +++ b/application/src/main/java/org/thingsboard/server/service/ws/notification/cmd/WsCmd.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.service.ws; +package org.thingsboard.server.service.ws.notification.cmd; public interface WsCmd { int getCmdId(); diff --git a/application/src/test/java/org/thingsboard/server/service/notification/NotificationsWsApiTest.java b/application/src/test/java/org/thingsboard/server/service/notification/NotificationApiTest.java similarity index 94% rename from application/src/test/java/org/thingsboard/server/service/notification/NotificationsWsApiTest.java rename to application/src/test/java/org/thingsboard/server/service/notification/NotificationApiTest.java index 6738b37a70..aca74f9a07 100644 --- a/application/src/test/java/org/thingsboard/server/service/notification/NotificationsWsApiTest.java +++ b/application/src/test/java/org/thingsboard/server/service/notification/NotificationApiTest.java @@ -35,7 +35,7 @@ import java.util.concurrent.TimeUnit; import static org.assertj.core.api.Assertions.assertThat; @DaoSqlTest -public class NotificationsWsApiTest extends AbstractControllerTest { +public class NotificationApiTest extends AbstractControllerTest { @Before public void beforeEach() throws Exception { @@ -154,8 +154,6 @@ public class NotificationsWsApiTest extends AbstractControllerTest { assertThat(getAnotherWsClient().getLastCountUpdate().getTotalUnreadCount()).isOne(); } - public void testReceivingUpdatesWhenSubscriptionAtAnotherInstance() {} - private void checkFullNotificationsUpdate(UnreadNotificationsUpdate notificationsUpdate, String... expectedNotifications) { assertThat(notificationsUpdate.getNotifications()).extracting(Notification::getText).containsOnly(expectedNotifications); @@ -189,19 +187,19 @@ public class NotificationsWsApiTest extends AbstractControllerTest { @Override protected TbTestWebSocketClient buildAndConnectWebSocketClient() throws URISyntaxException, InterruptedException { - NotificationsWebSocketClient wsClient = new NotificationsWebSocketClient(WS_URL + wsPort, token); + NotificationApiWsClient wsClient = new NotificationApiWsClient(WS_URL + wsPort, token); assertThat(wsClient.connectBlocking(TIMEOUT, TimeUnit.SECONDS)).isTrue(); return wsClient; } @Override - public NotificationsWebSocketClient getWsClient() { - return (NotificationsWebSocketClient) super.getWsClient(); + public NotificationApiWsClient getWsClient() { + return (NotificationApiWsClient) super.getWsClient(); } @Override - public NotificationsWebSocketClient getAnotherWsClient() { - return (NotificationsWebSocketClient) super.getAnotherWsClient(); + public NotificationApiWsClient getAnotherWsClient() { + return (NotificationApiWsClient) super.getAnotherWsClient(); } } diff --git a/application/src/test/java/org/thingsboard/server/service/notification/NotificationsWebSocketClient.java b/application/src/test/java/org/thingsboard/server/service/notification/NotificationApiWsClient.java similarity index 89% rename from application/src/test/java/org/thingsboard/server/service/notification/NotificationsWebSocketClient.java rename to application/src/test/java/org/thingsboard/server/service/notification/NotificationApiWsClient.java index 5da35b91fc..e5b9d4b887 100644 --- a/application/src/test/java/org/thingsboard/server/service/notification/NotificationsWebSocketClient.java +++ b/application/src/test/java/org/thingsboard/server/service/notification/NotificationApiWsClient.java @@ -21,7 +21,7 @@ import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.RandomUtils; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.controller.TbTestWebSocketClient; -import org.thingsboard.server.service.ws.notification.cmd.MarkNotificationAsReadCmd; +import org.thingsboard.server.service.ws.notification.cmd.MarkNotificationsAsReadCmd; import org.thingsboard.server.service.ws.notification.cmd.NotificationCmdsWrapper; import org.thingsboard.server.service.ws.notification.cmd.NotificationsCountSubCmd; import org.thingsboard.server.service.ws.notification.cmd.NotificationsSubCmd; @@ -31,17 +31,18 @@ import org.thingsboard.server.service.ws.telemetry.cmd.v2.CmdUpdateType; import java.net.URI; import java.net.URISyntaxException; +import java.util.Arrays; import java.util.UUID; @Slf4j -public class NotificationsWebSocketClient extends TbTestWebSocketClient { +public class NotificationApiWsClient extends TbTestWebSocketClient { @Getter private UnreadNotificationsUpdate lastDataUpdate; @Getter private UnreadNotificationsCountUpdate lastCountUpdate; - public NotificationsWebSocketClient(String wsUrl, String token) throws URISyntaxException { + public NotificationApiWsClient(String wsUrl, String token) throws URISyntaxException { super(new URI(wsUrl + "/api/ws/plugins/notifications?token=" + token)); } @@ -57,9 +58,9 @@ public class NotificationsWebSocketClient extends TbTestWebSocketClient { sendCmd(cmdsWrapper); } - public void markNotificationAsRead(UUID notificationId) { + public void markNotificationAsRead(UUID... notifications) { NotificationCmdsWrapper cmdsWrapper = new NotificationCmdsWrapper(); - cmdsWrapper.setMarkAsReadCmd(new MarkNotificationAsReadCmd(newCmdId(), notificationId)); + cmdsWrapper.setMarkAsReadCmd(new MarkNotificationsAsReadCmd(newCmdId(), Arrays.asList(notifications))); sendCmd(cmdsWrapper); } diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/notification/NotificationRequestService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/notification/NotificationRequestService.java index 4ded8cbd9f..7ae96ada79 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/notification/NotificationRequestService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/notification/NotificationRequestService.java @@ -15,7 +15,7 @@ */ package org.thingsboard.server.dao.notification; -import org.thingsboard.server.common.data.id.AlarmId; +import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.NotificationRequestId; import org.thingsboard.server.common.data.id.NotificationRuleId; import org.thingsboard.server.common.data.id.TenantId; @@ -33,7 +33,7 @@ public interface NotificationRequestService { PageData findNotificationRequestsByTenantId(TenantId tenantId, PageLink pageLink); - List findNotificationRequestsByRuleIdAndAlarmId(TenantId tenantId, NotificationRuleId ruleId, AlarmId alarmId); + List findNotificationRequestsByRuleIdAndOriginatorEntityId(TenantId tenantId, NotificationRuleId ruleId, EntityId originatorEntityId); void deleteNotificationRequestById(TenantId tenantId, NotificationRequestId id); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/notification/AlarmOriginatedNotificationInfo.java b/common/data/src/main/java/org/thingsboard/server/common/data/notification/AlarmOriginatedNotificationInfo.java new file mode 100644 index 0000000000..19d3dc99ad --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/notification/AlarmOriginatedNotificationInfo.java @@ -0,0 +1,46 @@ +/** + * 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.common.data.notification; + +import lombok.AllArgsConstructor; +import lombok.Builder; +import lombok.Data; +import lombok.EqualsAndHashCode; +import lombok.NoArgsConstructor; +import org.thingsboard.server.common.data.alarm.AlarmSeverity; +import org.thingsboard.server.common.data.alarm.AlarmStatus; +import org.thingsboard.server.common.data.id.AlarmId; +import org.thingsboard.server.common.data.id.EntityId; + +@Data +@EqualsAndHashCode(callSuper = true) +@NoArgsConstructor +@AllArgsConstructor +@Builder +public class AlarmOriginatedNotificationInfo extends NotificationInfo { + + private AlarmId alarmId; + private String alarmType; + private EntityId alarmOriginator; + private AlarmSeverity alarmSeverity; + private AlarmStatus alarmStatus; + + @Override + public NotificationOriginatorType getOriginatorType() { + return NotificationOriginatorType.ALARM; + } + +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/notification/Notification.java b/common/data/src/main/java/org/thingsboard/server/common/data/notification/Notification.java index 827bddc980..a39d858fcc 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/notification/Notification.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/notification/Notification.java @@ -38,7 +38,7 @@ public class Notification extends BaseData { private String text; private NotificationInfo info; private NotificationSeverity severity; + private NotificationOriginatorType originatorType; private NotificationStatus status; -// private UserId senderId; } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/notification/NotificationInfo.java b/common/data/src/main/java/org/thingsboard/server/common/data/notification/NotificationInfo.java index 298998f8d7..e7b2873d99 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/notification/NotificationInfo.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/notification/NotificationInfo.java @@ -15,6 +15,10 @@ */ package org.thingsboard.server.common.data.notification; +import com.fasterxml.jackson.annotation.JsonIgnoreProperties; +import com.fasterxml.jackson.annotation.JsonSubTypes; +import com.fasterxml.jackson.annotation.JsonSubTypes.Type; +import com.fasterxml.jackson.annotation.JsonTypeInfo; import lombok.AllArgsConstructor; import lombok.Builder; import lombok.Data; @@ -28,23 +32,18 @@ import org.thingsboard.server.common.data.validation.NoXss; @Data -@AllArgsConstructor -@NoArgsConstructor -@Builder -//@JsonIgnoreProperties(ignoreUnknown = true) -//@JsonTypeInfo(use = JsonTypeInfo.Id.NAME, property = "notificationType", visible = true, defaultImpl = NotificationInfo.class) -//@JsonSubTypes({ -// @Type(name = "ALARM", value = DeviceExportData.class), -//}) +@JsonIgnoreProperties(ignoreUnknown = true) +@JsonTypeInfo(use = JsonTypeInfo.Id.NAME, property = "originatorType", defaultImpl = NotificationInfo.class) +@JsonSubTypes({ + @Type(name = "ALARM", value = AlarmOriginatedNotificationInfo.class), +}) public class NotificationInfo { - @NoXss + private String description; private DashboardId dashboardId; - private AlarmId alarmId; - private String alarmType; - private EntityId alarmOriginator; - private AlarmSeverity alarmSeverity; - private AlarmStatus alarmStatus; + public NotificationOriginatorType getOriginatorType() { + return null; + } } diff --git a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineNotificationService.java b/common/data/src/main/java/org/thingsboard/server/common/data/notification/NotificationOriginatorType.java similarity index 64% rename from rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineNotificationService.java rename to common/data/src/main/java/org/thingsboard/server/common/data/notification/NotificationOriginatorType.java index e7b87cf25d..44aa879a00 100644 --- a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineNotificationService.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/notification/NotificationOriginatorType.java @@ -13,13 +13,10 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.rule.engine.api; - -import org.thingsboard.server.common.data.id.TenantId; -import org.thingsboard.server.common.data.notification.NotificationRequest; - -public interface RuleEngineNotificationService { - - NotificationRequest processNotificationRequest(TenantId tenantId, NotificationRequest notificationRequest); +package org.thingsboard.server.common.data.notification; +public enum NotificationOriginatorType { + USER, + ALARM, + RULE_NODE } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/notification/NotificationRequest.java b/common/data/src/main/java/org/thingsboard/server/common/data/notification/NotificationRequest.java index af90eab7ee..f962db97e4 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/notification/NotificationRequest.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/notification/NotificationRequest.java @@ -16,7 +16,6 @@ package org.thingsboard.server.common.data.notification; import com.fasterxml.jackson.annotation.JsonIgnore; -import com.fasterxml.jackson.annotation.JsonProperty; import lombok.AllArgsConstructor; import lombok.Builder; import lombok.Data; @@ -25,7 +24,7 @@ import lombok.NoArgsConstructor; import org.thingsboard.server.common.data.BaseData; import org.thingsboard.server.common.data.HasName; import org.thingsboard.server.common.data.HasTenantId; -import org.thingsboard.server.common.data.id.AlarmId; +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; @@ -46,26 +45,25 @@ public class NotificationRequest extends BaseData impleme private TenantId tenantId; @NotNull(message = "Target is not specified") private NotificationTargetId targetId; + @NoXss - private String notificationReason; // "Alarm", "Scheduled event". "General" by default - // @NoXss + private String notificationReason; @NotBlank(message = "Notification text template is missing") - private String textTemplate; + private String textTemplate; // fixme: xss @Valid private NotificationInfo notificationInfo; private NotificationSeverity notificationSeverity; - private NotificationRequestConfig additionalConfig; - @JsonProperty(access = JsonProperty.Access.READ_ONLY) - private NotificationRequestStatus status; - @JsonIgnore + private NotificationOriginatorType originatorType; + private EntityId originatorEntityId; // userId, alarmId or tenantId private NotificationRuleId ruleId; // maybe move to child class - @JsonIgnore - private AlarmId alarmId; + + private NotificationRequestConfig additionalConfig; + private NotificationRequestStatus status; public static final String GENERAL_NOTIFICATION_REASON = "General"; - public static final String ALARM_NOTIFICATION_REASON = "Alarm"; + @JsonIgnore @Override public String getName() { return notificationReason; diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java b/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java index 27b17164ea..9d477059d6 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java @@ -660,6 +660,7 @@ public class ModelConstants { public static final String NOTIFICATION_TEXT_PROPERTY = "text"; public static final String NOTIFICATION_INFO_PROPERTY = "info"; public static final String NOTIFICATION_SEVERITY_PROPERTY = "severity"; + public static final String NOTIFICATION_ORIGINATOR_TYPE_PROPERTY = "originator_type"; public static final String NOTIFICATION_STATUS_PROPERTY = "status"; public static final String NOTIFICATION_REQUEST_TABLE_NAME = "notification_request"; @@ -668,10 +669,12 @@ public class ModelConstants { public static final String NOTIFICATION_REQUEST_NOTIFICATION_REASON_PROPERTY = "notification_reason"; public static final String NOTIFICATION_REQUEST_NOTIFICATION_INFO_PROPERTY = "notification_info"; public static final String NOTIFICATION_REQUEST_NOTIFICATION_SEVERITY_PROPERTY = "notification_severity"; + public static final String NOTIFICATION_REQUEST_ORIGINATOR_TYPE_PROPERTY = "originator_type"; + public static final String NOTIFICATION_REQUEST_ORIGINATOR_ENTITY_ID_PROPERTY = "originator_entity_id"; + public static final String NOTIFICATION_REQUEST_ORIGINATOR_ENTITY_TYPE_PROPERTY = "originator_entity_type"; public static final String NOTIFICATION_REQUEST_ADDITIONAL_CONFIG_PROPERTY = "additional_config"; public static final String NOTIFICATION_REQUEST_STATUS_PROPERTY = "status"; public static final String NOTIFICATION_REQUEST_RULE_ID_PROPERTY = "rule_id"; - public static final String NOTIFICATION_REQUEST_ALARM_ID_PROPERTY = "alarm_id"; public static final String NOTIFICATION_RULE_TABLE_NAME = "notification_rule"; // ... diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/sql/NotificationEntity.java b/dao/src/main/java/org/thingsboard/server/dao/model/sql/NotificationEntity.java index c79ba6f30b..c94ec80016 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/sql/NotificationEntity.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/sql/NotificationEntity.java @@ -26,6 +26,7 @@ import org.thingsboard.server.common.data.id.NotificationRequestId; import org.thingsboard.server.common.data.id.UserId; 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 org.thingsboard.server.common.data.notification.NotificationSeverity; import org.thingsboard.server.common.data.notification.NotificationStatus; import org.thingsboard.server.dao.model.BaseSqlEntity; @@ -66,6 +67,10 @@ public class NotificationEntity extends BaseSqlEntity { @Column(name = ModelConstants.NOTIFICATION_SEVERITY_PROPERTY) private NotificationSeverity severity; + @Enumerated(EnumType.STRING) + @Column(name = ModelConstants.NOTIFICATION_ORIGINATOR_TYPE_PROPERTY) + private NotificationOriginatorType originatorType; + @Enumerated(EnumType.STRING) @Column(name = ModelConstants.NOTIFICATION_STATUS_PROPERTY) private NotificationStatus status; @@ -83,6 +88,7 @@ public class NotificationEntity extends BaseSqlEntity { setInfo(JacksonUtil.valueToTree(notification.getInfo())); } setSeverity(notification.getSeverity()); + setOriginatorType(notification.getOriginatorType()); setStatus(notification.getStatus()); } @@ -99,6 +105,7 @@ public class NotificationEntity extends BaseSqlEntity { notification.setInfo(JacksonUtil.treeToValue(info, NotificationInfo.class)); } notification.setSeverity(severity); + notification.setOriginatorType(originatorType); notification.setStatus(status); return notification; } diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/sql/NotificationRequestEntity.java b/dao/src/main/java/org/thingsboard/server/dao/model/sql/NotificationRequestEntity.java index 3c7b525628..52b5aee2b9 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/sql/NotificationRequestEntity.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/sql/NotificationRequestEntity.java @@ -21,13 +21,16 @@ import lombok.EqualsAndHashCode; import org.hibernate.annotations.Type; import org.hibernate.annotations.TypeDef; import org.thingsboard.common.util.JacksonUtil; +import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.id.AlarmId; +import org.thingsboard.server.common.data.id.EntityIdFactory; 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; import org.thingsboard.server.common.data.id.UserId; 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; @@ -70,6 +73,20 @@ public class NotificationRequestEntity extends BaseSqlEntity findNotificationRequestsByRuleIdAndAlarmId(TenantId tenantId, NotificationRuleId ruleId, AlarmId alarmId) { - return notificationRequestDao.findByRuleIdAndAlarmId(tenantId, ruleId, alarmId); + public List findNotificationRequestsByRuleIdAndOriginatorEntityId(TenantId tenantId, NotificationRuleId ruleId, EntityId originatorEntityId) { + return notificationRequestDao.findByRuleIdAndOriginatorEntityId(tenantId, ruleId, originatorEntityId); } // ON DELETE CASCADE is used: notifications for request are deleted as well diff --git a/dao/src/main/java/org/thingsboard/server/dao/notification/NotificationRequestDao.java b/dao/src/main/java/org/thingsboard/server/dao/notification/NotificationRequestDao.java index 76778c64b5..2036105ffb 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/notification/NotificationRequestDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/notification/NotificationRequestDao.java @@ -15,7 +15,7 @@ */ package org.thingsboard.server.dao.notification; -import org.thingsboard.server.common.data.id.AlarmId; +import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.NotificationRuleId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.notification.NotificationRequest; @@ -30,7 +30,7 @@ public interface NotificationRequestDao extends Dao { PageData findByTenantIdAndPageLink(TenantId tenantId, PageLink pageLink); - List findByRuleIdAndAlarmId(TenantId tenantId, NotificationRuleId ruleId, AlarmId alarmId); + List findByRuleIdAndOriginatorEntityId(TenantId tenantId, NotificationRuleId ruleId, EntityId originatorEntityId); PageData findAllByStatus(NotificationRequestStatus status, PageLink pageLink); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/notification/JpaNotificationRequestDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/notification/JpaNotificationRequestDao.java index 7684ed5865..b58fc198a7 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/notification/JpaNotificationRequestDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/notification/JpaNotificationRequestDao.java @@ -19,7 +19,7 @@ import com.google.common.base.Strings; import lombok.RequiredArgsConstructor; import org.springframework.data.jpa.repository.JpaRepository; import org.springframework.stereotype.Component; -import org.thingsboard.server.common.data.id.AlarmId; +import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.NotificationRuleId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.notification.NotificationRequest; @@ -49,8 +49,8 @@ public class JpaNotificationRequestDao extends JpaAbstractDao findByRuleIdAndAlarmId(TenantId tenantId, NotificationRuleId ruleId, AlarmId alarmId) { - return DaoUtil.convertDataList(notificationRequestRepository.findAllByRuleIdAndAlarmId(ruleId.getId(), alarmId.getId())); + public List findByRuleIdAndOriginatorEntityId(TenantId tenantId, NotificationRuleId ruleId, EntityId originatorEntityId) { + return DaoUtil.convertDataList(notificationRequestRepository.findAllByRuleIdAndOriginatorEntityTypeAndOriginatorEntityId(ruleId.getId(), originatorEntityId.getEntityType(), originatorEntityId.getId())); } @Override diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/notification/NotificationRequestRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/notification/NotificationRequestRepository.java index f892767a83..40ae696979 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/notification/NotificationRequestRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/notification/NotificationRequestRepository.java @@ -21,6 +21,7 @@ import org.springframework.data.jpa.repository.JpaRepository; import org.springframework.data.jpa.repository.Query; import org.springframework.data.repository.query.Param; import org.springframework.stereotype.Repository; +import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.notification.NotificationRequestStatus; import org.thingsboard.server.dao.model.sql.NotificationRequestEntity; @@ -36,7 +37,7 @@ public interface NotificationRequestRepository extends JpaRepository findByTenantIdAndSearchText(@Param("tenantId") UUID tenantId, @Param("searchText") String searchText, Pageable pageable); - List findAllByRuleIdAndAlarmId(UUID ruleId, UUID alarmId); + List findAllByRuleIdAndOriginatorEntityTypeAndOriginatorEntityId(UUID ruleId, EntityType originatorEntityType, UUID originatorEntityId); Page findAllByStatus(NotificationRequestStatus status, Pageable pageable); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/query/AlarmDataAdapter.java b/dao/src/main/java/org/thingsboard/server/dao/sql/query/AlarmDataAdapter.java index ff156eb51c..b825d36892 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/query/AlarmDataAdapter.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/query/AlarmDataAdapter.java @@ -27,6 +27,7 @@ import org.thingsboard.server.common.data.id.AlarmId; import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityIdFactory; +import org.thingsboard.server.common.data.id.NotificationRuleId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.query.AlarmData; @@ -101,6 +102,9 @@ public class AlarmDataAdapter { } else { alarm.setPropagateRelationTypes(Collections.emptyList()); } + if (row.get(ModelConstants.ALARM_NOTIFICATION_RULE_ID) != null) { + alarm.setNotificationRuleId(new NotificationRuleId((UUID) row.get(ModelConstants.ALARM_NOTIFICATION_RULE_ID))); + } UUID entityUuid = (UUID) row.get(ModelConstants.ENTITY_ID_COLUMN); EntityId entityId = entityIdMap.get(entityUuid); Object originatorNameObj = row.get(ModelConstants.ALARM_ORIGINATOR_NAME_PROPERTY); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/query/DefaultAlarmQueryRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/query/DefaultAlarmQueryRepository.java index 69cb4a2d86..94eebc7f6b 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/query/DefaultAlarmQueryRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/query/DefaultAlarmQueryRepository.java @@ -104,6 +104,7 @@ public class DefaultAlarmQueryRepository implements AlarmQueryRepository { " a.tenant_id as tenant_id, " + " a.customer_id as customer_id, " + " a.propagate_relation_types as propagate_relation_types, " + + " a.notification_rule_id as notification_rule_id, " + " a.type as type," + SELECT_ORIGINATOR_NAME + ", "; private static final String JOIN_ENTITY_ALARMS = "inner join entity_alarm ea on a.id = ea.alarm_id"; diff --git a/dao/src/main/resources/sql/schema-entities-idx.sql b/dao/src/main/resources/sql/schema-entities-idx.sql index 019a8b1771..0b95809b91 100644 --- a/dao/src/main/resources/sql/schema-entities-idx.sql +++ b/dao/src/main/resources/sql/schema-entities-idx.sql @@ -74,12 +74,12 @@ CREATE INDEX IF NOT EXISTS idx_rule_node_type ON rule_node(type); CREATE INDEX IF NOT EXISTS idx_api_usage_state_entity_id ON api_usage_state(entity_id); -CREATE INDEX IF NOT EXISTS idx_notification_target_tenant_id_and_created_time ON notification_target(tenant_id, created_time DESC); +CREATE INDEX IF NOT EXISTS idx_notification_target_tenant_id_created_time ON notification_target(tenant_id, created_time DESC); -CREATE INDEX IF NOT EXISTS idx_notification_request_tenant_id_and_created_time ON notification_request(tenant_id, created_time DESC); +CREATE INDEX IF NOT EXISTS idx_notification_request_tenant_id_originator_type_created_time ON notification_request(tenant_id, originator_type, created_time DESC); CREATE INDEX IF NOT EXISTS idx_notification_id ON notification(id); -CREATE INDEX IF NOT EXISTS idx_notification_notification_request_id ON notification(request_id); +CREATE INDEX IF NOT EXISTS idx_notification_recipient_id_created_time ON notification(recipient_id, created_time DESC); -CREATE INDEX IF NOT EXISTS idx_notification_recipient_id_and_created_time ON notification(recipient_id, created_time DESC); +CREATE INDEX IF NOT EXISTS idx_notification_notification_request_id ON notification(request_id); diff --git a/dao/src/main/resources/sql/schema-entities.sql b/dao/src/main/resources/sql/schema-entities.sql index 1ef732b8b6..885608c7f5 100644 --- a/dao/src/main/resources/sql/schema-entities.sql +++ b/dao/src/main/resources/sql/schema-entities.sql @@ -58,7 +58,8 @@ CREATE TABLE IF NOT EXISTS alarm ( propagate_relation_types varchar, type varchar(255), propagate_to_owner boolean, - propagate_to_tenant boolean + propagate_to_tenant boolean, + notification_rule_id uuid ); CREATE TABLE IF NOT EXISTS entity_alarm ( @@ -799,10 +800,12 @@ CREATE TABLE IF NOT EXISTS notification_request ( text_template VARCHAR NOT NULL, notification_info VARCHAR(1000), notification_severity VARCHAR(32), - additional_config VARCHAR(1000), - status VARCHAR(32), + 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), - alarm_id UUID + additional_config VARCHAR(1000), + status VARCHAR(32) ); CREATE TABLE IF NOT EXISTS notification ( @@ -815,5 +818,6 @@ CREATE TABLE IF NOT EXISTS notification ( text VARCHAR NOT NULL, info VARCHAR(1000), severity VARCHAR(32), + originator_type VARCHAR(32) NOT NULL, status VARCHAR(32) ) PARTITION BY RANGE (created_time); diff --git a/application/src/main/java/org/thingsboard/server/service/notification/NotificationSubscriptionService.java b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/NotificationManager.java similarity index 92% rename from application/src/main/java/org/thingsboard/server/service/notification/NotificationSubscriptionService.java rename to rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/NotificationManager.java index 902ad39470..ac336e9fe4 100644 --- a/application/src/main/java/org/thingsboard/server/service/notification/NotificationSubscriptionService.java +++ b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/NotificationManager.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.service.notification; +package org.thingsboard.rule.engine.api; import org.thingsboard.server.common.data.id.NotificationId; import org.thingsboard.server.common.data.id.NotificationRequestId; @@ -21,7 +21,7 @@ import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.UserId; import org.thingsboard.server.common.data.notification.NotificationRequest; -public interface NotificationSubscriptionService { +public interface NotificationManager { NotificationRequest processNotificationRequest(TenantId tenantId, NotificationRequest notificationRequest); diff --git a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java index 7ca31da24e..ede2f716e5 100644 --- a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java +++ b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java @@ -267,7 +267,7 @@ public interface TbContext { SmsSenderFactory getSmsSenderFactory(); - RuleEngineNotificationService getNotificationService(); + NotificationManager getNotificationManager(); /** * Creates JS Script Engine diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbNotificationNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbNotificationNode.java index 06fc0b110d..f87debd4dc 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbNotificationNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbNotificationNode.java @@ -23,6 +23,7 @@ import org.thingsboard.rule.engine.api.TbNodeConfiguration; import org.thingsboard.rule.engine.api.TbNodeException; import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.server.common.data.id.NotificationTargetId; +import org.thingsboard.server.common.data.notification.NotificationOriginatorType; import org.thingsboard.server.common.data.notification.NotificationRequest; import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.msg.TbMsg; @@ -58,9 +59,11 @@ public class TbNotificationNode implements TbNode { .notificationReason(config.getNotificationReason()) .textTemplate(TbNodeUtils.processPattern(config.getNotificationTextTemplate(), msg)) .notificationSeverity(config.getNotificationSeverity()) + .originatorType(NotificationOriginatorType.RULE_NODE) + .originatorEntityId(ctx.getTenantId()) .build(); withCallback(ctx.getDbCallbackExecutor().executeAsync(() -> { - return ctx.getNotificationService().processNotificationRequest(ctx.getTenantId(), notificationRequest); + return ctx.getNotificationManager().processNotificationRequest(ctx.getTenantId(), notificationRequest); }), r -> { TbMsgMetaData msgMetaData = msg.getMetaData().copy();