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 a4f44b589d..5c621477a1 100644 --- a/application/src/main/java/org/thingsboard/server/controller/NotificationController.java +++ b/application/src/main/java/org/thingsboard/server/controller/NotificationController.java @@ -37,6 +37,7 @@ import org.thingsboard.server.common.data.notification.Notification; import org.thingsboard.server.common.data.notification.NotificationRequest; 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; @@ -54,6 +55,7 @@ import java.util.UUID; public class NotificationController extends BaseController { private final NotificationService notificationService; + private final NotificationRequestService notificationRequestService; private final NotificationSubscriptionService notificationSubscriptionService; @GetMapping("/notifications") @@ -102,7 +104,7 @@ public class NotificationController extends BaseController { public NotificationRequest getNotificationRequestById(@PathVariable UUID id, @AuthenticationPrincipal SecurityUser user) { NotificationRequestId notificationRequestId = new NotificationRequestId(id); - return notificationService.findNotificationRequestById(user.getTenantId(), notificationRequestId); + return notificationRequestService.findNotificationRequestById(user.getTenantId(), notificationRequestId); } @GetMapping("/notification/requests") @@ -114,14 +116,14 @@ public class NotificationController extends BaseController { @RequestParam(required = false) String sortOrder, @AuthenticationPrincipal SecurityUser user) throws ThingsboardException { PageLink pageLink = createPageLink(pageSize, page, textSearch, sortProperty, sortOrder); - return notificationService.findNotificationRequestsByTenantId(user.getTenantId(), pageLink); + return notificationRequestService.findNotificationRequestsByTenantId(user.getTenantId(), pageLink); } @DeleteMapping("/notification/request/{id}") public void deleteNotificationRequest(@PathVariable UUID id, @AuthenticationPrincipal SecurityUser user) throws ThingsboardException { NotificationRequestId notificationRequestId = new NotificationRequestId(id); - NotificationRequest notificationRequest = notificationService.findNotificationRequestById(user.getTenantId(), notificationRequestId); + NotificationRequest notificationRequest = notificationRequestService.findNotificationRequestById(user.getTenantId(), notificationRequestId); accessControlService.checkPermission(user, Resource.NOTIFICATION_REQUEST, Operation.DELETE, notificationRequestId, notificationRequest); try { notificationSubscriptionService.deleteNotificationRequest(user.getTenantId(), notificationRequestId); 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 7c459b2e80..a36b5cc8c2 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 @@ -30,8 +30,8 @@ import org.thingsboard.server.common.data.notification.NotificationRequestStatus import org.thingsboard.server.common.data.notification.NotificationSeverity; import org.thingsboard.server.common.data.notification.rule.NonConfirmedNotificationEscalation; import org.thingsboard.server.common.data.notification.rule.NotificationRule; +import org.thingsboard.server.dao.notification.NotificationRequestService; import org.thingsboard.server.dao.notification.NotificationRuleService; -import org.thingsboard.server.dao.notification.NotificationService; import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.executors.DbCallbackExecutorService; @@ -43,7 +43,7 @@ import java.util.List; public class DefaultNotificationRuleProcessingService implements NotificationRuleProcessingService { private final NotificationRuleService notificationRuleService; - private final NotificationService notificationService; + private final NotificationRequestService notificationRequestService; private final NotificationSubscriptionService notificationSubscriptionService; private final DbCallbackExecutorService dbCallbackExecutorService; @@ -65,7 +65,7 @@ public class DefaultNotificationRuleProcessingService implements NotificationRul } private void onAlarmUpdate(TenantId tenantId, NotificationRuleId notificationRuleId, Alarm alarm) { - List notificationRequests = notificationService.findNotificationRequestsByRuleIdAndAlarmId(tenantId, notificationRuleId, alarm.getId()); + List notificationRequests = notificationRequestService.findNotificationRequestsByRuleIdAndAlarmId(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 @@ -80,7 +80,7 @@ public class DefaultNotificationRuleProcessingService implements NotificationRul for (NotificationRequest notificationRequest : notificationRequests) { if (notificationRequest.getStatus() == NotificationRequestStatus.SCHEDULED) { // using regular service due to no need to send an update to subscription manager - notificationService.deleteNotificationRequestById(tenantId, notificationRequest.getId()); + notificationRequestService.deleteNotificationRequestById(tenantId, notificationRequest.getId()); } else { notificationSubscriptionService.deleteNotificationRequest(tenantId, notificationRequest.getId()); // todo: or should we mark already sent notifications as read? 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 new file mode 100644 index 0000000000..903444529b --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationSchedulerService.java @@ -0,0 +1,129 @@ +/** + * 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.util.concurrent.ListenableFuture; +import com.google.common.util.concurrent.ListenableScheduledFuture; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Service; +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.msg.queue.ServiceType; +import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; +import org.thingsboard.server.dao.notification.NotificationRequestService; +import org.thingsboard.server.queue.util.TbCoreComponent; +import org.thingsboard.server.service.partition.AbstractPartitionBasedService; + +import javax.annotation.PostConstruct; +import java.util.Collections; +import java.util.List; +import java.util.Map; +import java.util.Optional; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ScheduledFuture; +import java.util.concurrent.TimeUnit; + +@TbCoreComponent +@Service +@RequiredArgsConstructor +@Slf4j +@SuppressWarnings("UnstableApiUsage") +public class DefaultNotificationSchedulerService extends AbstractPartitionBasedService implements NotificationSchedulerService { + + private final NotificationSubscriptionService notificationSubscriptionService; + private final NotificationRequestService notificationRequestService; + + private final Map> scheduledNotificationRequests = new ConcurrentHashMap<>(); + + @PostConstruct + public void init() { + super.init(); + } + + @Override + protected Map>> onAddedPartitions(Set addedPartitions) { + PageDataIterable notificationRequests = new PageDataIterable<>(pageLink -> { + return notificationRequestService.findScheduledNotificationRequests(pageLink); + }, 1000); + for (NotificationRequest notificationRequest : notificationRequests) { + TopicPartitionInfo requestPartition = partitionService.resolve(ServiceType.TB_CORE, notificationRequest.getTenantId(), notificationRequest.getId()); + if (addedPartitions.contains(requestPartition)) { + partitionedEntities.computeIfAbsent(requestPartition, k -> ConcurrentHashMap.newKeySet()).add(notificationRequest.getId()); + if (!scheduledNotificationRequests.containsKey(notificationRequest.getId())) { + scheduleNotificationRequest(notificationRequest.getTenantId(), notificationRequest, notificationRequest.getCreatedTime()); + } + } + } + return Collections.emptyMap(); + } + + @Override + public void scheduleNotificationRequest(TenantId tenantId, NotificationRequestId notificationRequestId, long requestTs) { + NotificationRequest notificationRequest = notificationRequestService.findNotificationRequestById(tenantId, notificationRequestId); + scheduleNotificationRequest(tenantId, notificationRequest, requestTs); + } + + private void scheduleNotificationRequest(TenantId tenantId, NotificationRequest request, long requestTs) { + int delayInMinutes = Optional.ofNullable(request) + .map(NotificationRequest::getAdditionalConfig) + .map(NotificationRequestConfig::getSendingDelayInMinutes) + .orElse(0); + if (delayInMinutes <= 0) return; // todo: think about: if server was down for some time and delayMs will be negative - need to send these requests as well (but when the value is within some range) + long delayMs = TimeUnit.MINUTES.toMillis(delayInMinutes) - (System.currentTimeMillis() - requestTs); + + ListenableScheduledFuture scheduledTask = scheduledExecutor.schedule(() -> { + NotificationRequest notificationRequest = notificationRequestService.findNotificationRequestById(tenantId, request.getId()); + if (notificationRequest == null) return; + + notificationSubscriptionService.processNotificationRequest(tenantId, notificationRequest); + scheduledNotificationRequests.remove(notificationRequest.getId()); + }, delayMs, TimeUnit.MILLISECONDS); + scheduledNotificationRequests.put(request.getId(), scheduledTask); + } + + @Override + public void onNotificationRequestDeleted(TenantId tenantId, NotificationRequestId notificationRequestId) { + removeAndCancel(notificationRequestId); + } + + @Override + protected void cleanupEntityOnPartitionRemoval(NotificationRequestId notificationRequestId) { + removeAndCancel(notificationRequestId); + } + + private void removeAndCancel(NotificationRequestId notificationRequestId) { + ScheduledFuture scheduledTask = scheduledNotificationRequests.remove(notificationRequestId); + if (scheduledTask != null) { + scheduledTask.cancel(false); + } + } + + @Override + protected String getServiceName() { + return "Notifications scheduler"; + } + + @Override + protected String getSchedulerExecutorName() { + return "notifications-scheduler"; + } + +} diff --git a/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationSubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationSubscriptionService.java index 74fb400fb2..8c786f33c9 100644 --- a/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationSubscriptionService.java +++ b/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationSubscriptionService.java @@ -35,6 +35,7 @@ import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.dao.DaoUtil; +import org.thingsboard.server.dao.notification.NotificationRequestService; import org.thingsboard.server.dao.notification.NotificationService; import org.thingsboard.server.dao.notification.NotificationTargetService; import org.thingsboard.server.gen.transport.TransportProtos; @@ -56,17 +57,20 @@ import java.util.UUID; public class DefaultNotificationSubscriptionService extends AbstractSubscriptionService implements NotificationSubscriptionService, RuleEngineNotificationService { private final NotificationTargetService notificationTargetService; + private final NotificationRequestService notificationRequestService; private final NotificationService notificationService; private final DbCallbackExecutorService dbCallbackExecutorService; private final NotificationsTopicService notificationsTopicService; public DefaultNotificationSubscriptionService(TbClusterService clusterService, PartitionService partitionService, NotificationTargetService notificationTargetService, + NotificationRequestService notificationRequestService, NotificationService notificationService, DbCallbackExecutorService dbCallbackExecutorService, NotificationsTopicService notificationsTopicService) { super(clusterService, partitionService); this.notificationTargetService = notificationTargetService; + this.notificationRequestService = notificationRequestService; this.notificationService = notificationService; this.dbCallbackExecutorService = dbCallbackExecutorService; this.notificationsTopicService = notificationsTopicService; @@ -77,18 +81,18 @@ 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) { + if (config.getSendingDelayInMinutes() > 0 && notificationRequest.getId() == null) { notificationRequest.setStatus(NotificationRequestStatus.SCHEDULED); - - // todo: delayed sending; check all delayed notification requests on start up, schedule send + NotificationRequest savedNotificationRequest = notificationRequestService.saveNotificationRequest(tenantId, notificationRequest); + forwardToNotificationSchedulerService(tenantId, savedNotificationRequest.getId(), false); + return savedNotificationRequest; } } - if (notificationRequest.getStatus() == null) { - notificationRequest.setStatus(NotificationRequestStatus.PROCESSED); - } - NotificationRequest savedNotificationRequest = notificationService.saveNotificationRequest(tenantId, notificationRequest); + notificationRequest.setStatus(NotificationRequestStatus.PROCESSED); + NotificationRequest savedNotificationRequest = notificationRequestService.saveNotificationRequest(tenantId, notificationRequest); DaoUtil.processBatches(pageLink -> { return notificationTargetService.findRecipientsForNotificationTarget(tenantId, notificationRequest.getTargetId(), pageLink); @@ -109,6 +113,20 @@ public class DefaultNotificationSubscriptionService extends AbstractSubscription return savedNotificationRequest; } + private void forwardToNotificationSchedulerService(TenantId tenantId, NotificationRequestId notificationRequestId, boolean deleted) { + 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); + TransportProtos.ToCoreMsg toCoreMsg = TransportProtos.ToCoreMsg.newBuilder() + .setNotificationSchedulerServiceMsg(msg) + .build(); + clusterService.pushMsgToCore(tenantId, notificationRequestId, toCoreMsg, null); + } + @Override public void markNotificationAsRead(TenantId tenantId, UserId recipientId, NotificationId notificationId) { boolean updated = notificationService.markNotificationAsRead(tenantId, recipientId, notificationId); @@ -122,17 +140,18 @@ public class DefaultNotificationSubscriptionService extends AbstractSubscription @Override public void deleteNotificationRequest(TenantId tenantId, NotificationRequestId notificationRequestId) { log.debug("Deleting notification request {}", notificationRequestId); - notificationService.deleteNotificationRequestById(tenantId, notificationRequestId); + notificationRequestService.deleteNotificationRequestById(tenantId, notificationRequestId); onNotificationRequestUpdate(tenantId, NotificationRequestUpdate.builder() .notificationRequestId(notificationRequestId) .deleted(true) .build()); + forwardToNotificationSchedulerService(tenantId, notificationRequestId, true); } @Override public void updateNotificationRequest(TenantId tenantId, NotificationRequest notificationRequest) { log.debug("Updating notification request {}", notificationRequest.getId()); - notificationService.saveNotificationRequest(tenantId, notificationRequest); + notificationRequestService.saveNotificationRequest(tenantId, notificationRequest); notificationService.updateNotificationsInfosByRequestId(tenantId, notificationRequest.getId(), notificationRequest.getNotificationInfo()); onNotificationRequestUpdate(tenantId, NotificationRequestUpdate.builder() .notificationRequestId(notificationRequest.getId()) diff --git a/application/src/main/java/org/thingsboard/server/service/notification/NotificationSchedulerService.java b/application/src/main/java/org/thingsboard/server/service/notification/NotificationSchedulerService.java new file mode 100644 index 0000000000..ce018f171c --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/notification/NotificationSchedulerService.java @@ -0,0 +1,27 @@ +/** + * 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.thingsboard.server.common.data.id.NotificationRequestId; +import org.thingsboard.server.common.data.id.TenantId; + +public interface NotificationSchedulerService { + + void scheduleNotificationRequest(TenantId tenantId, NotificationRequestId notificationRequestId, long requestTs); + + void onNotificationRequestDeleted(TenantId tenantId, NotificationRequestId notificationRequestId); + +} diff --git a/application/src/main/java/org/thingsboard/server/service/partition/AbstractPartitionBasedService.java b/application/src/main/java/org/thingsboard/server/service/partition/AbstractPartitionBasedService.java index 305696662b..489b1dcb48 100644 --- a/application/src/main/java/org/thingsboard/server/service/partition/AbstractPartitionBasedService.java +++ b/application/src/main/java/org/thingsboard/server/service/partition/AbstractPartitionBasedService.java @@ -20,16 +20,17 @@ import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListeningScheduledExecutorService; import com.google.common.util.concurrent.MoreExecutors; import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.annotation.Autowired; import org.thingsboard.common.util.DonAsynchron; import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; +import org.thingsboard.server.queue.discovery.PartitionService; import org.thingsboard.server.queue.discovery.TbApplicationEventListener; import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent; import java.util.ArrayList; -import java.util.Collection; import java.util.HashSet; import java.util.List; import java.util.Map; @@ -48,6 +49,8 @@ public abstract class AbstractPartitionBasedService extends protected final ConcurrentMap>> partitionedFetchTasks = new ConcurrentHashMap<>(); final Queue> subscribeQueue = new ConcurrentLinkedQueue<>(); + @Autowired + protected PartitionService partitionService; protected ListeningScheduledExecutorService scheduledExecutor; abstract protected String getServiceName(); diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java index cf93ea7f1d..60a9567488 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java @@ -31,8 +31,6 @@ import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.NotificationRequestId; import org.thingsboard.server.common.data.id.TenantId; 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.rpc.RpcError; import org.thingsboard.server.common.msg.MsgType; import org.thingsboard.server.common.msg.TbActorMsg; @@ -69,6 +67,7 @@ import org.thingsboard.server.queue.util.AfterStartUp; import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.apiusage.TbApiUsageStateService; import org.thingsboard.server.service.edge.EdgeNotificationService; +import org.thingsboard.server.service.notification.NotificationSchedulerService; import org.thingsboard.server.service.ota.OtaPackageStateService; import org.thingsboard.server.service.profile.TbAssetProfileCache; import org.thingsboard.server.service.profile.TbDeviceProfileCache; @@ -127,6 +126,7 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService> usageStatsConsumer; private final TbQueueConsumer> firmwareStatesConsumer; @@ -151,7 +151,8 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService findNotificationRequestsByTenantId(TenantId tenantId, PageLink pageLink); + + List findNotificationRequestsByRuleIdAndAlarmId(TenantId tenantId, NotificationRuleId ruleId, AlarmId alarmId); + + void deleteNotificationRequestById(TenantId tenantId, NotificationRequestId id); + + PageData findScheduledNotificationRequests(PageLink pageLink); + +} diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/notification/NotificationService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/notification/NotificationService.java index 6b7f6267db..498e3dfa3f 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/notification/NotificationService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/notification/NotificationService.java @@ -15,33 +15,17 @@ */ package org.thingsboard.server.dao.notification; -import org.thingsboard.server.common.data.id.AlarmId; import org.thingsboard.server.common.data.id.NotificationId; import org.thingsboard.server.common.data.id.NotificationRequestId; -import org.thingsboard.server.common.data.id.NotificationRuleId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.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.NotificationRequest; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageLink; -import java.util.List; - public interface NotificationService { - NotificationRequest saveNotificationRequest(TenantId tenantId, NotificationRequest notificationRequest); - - NotificationRequest findNotificationRequestById(TenantId tenantId, NotificationRequestId id); - - PageData findNotificationRequestsByTenantId(TenantId tenantId, PageLink pageLink); - - List findNotificationRequestsByRuleIdAndAlarmId(TenantId tenantId, NotificationRuleId ruleId, AlarmId alarmId); - - void deleteNotificationRequestById(TenantId tenantId, NotificationRequestId id); - - Notification saveNotification(TenantId tenantId, Notification notification); Notification findNotificationById(TenantId tenantId, NotificationId notificationId); diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/BaseSqlEntity.java b/dao/src/main/java/org/thingsboard/server/dao/model/BaseSqlEntity.java index 5256eee7db..d639192917 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/BaseSqlEntity.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/BaseSqlEntity.java @@ -86,8 +86,8 @@ public abstract class BaseSqlEntity implements BaseEntity { } } - protected T fromJson(JsonNode json) { - return JacksonUtil.convertValue(json, new TypeReference() {}); + protected T fromJson(JsonNode json, Class type) { + return JacksonUtil.convertValue(json, type); } } 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 cd0bbd435a..3c7b525628 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 @@ -110,9 +110,9 @@ public class NotificationRequestEntity extends BaseSqlEntity findNotificationRequestsByTenantId(TenantId tenantId, PageLink pageLink) { + return notificationRequestDao.findByTenantIdAndPageLink(tenantId, pageLink); + } + + @Override + public List findNotificationRequestsByRuleIdAndAlarmId(TenantId tenantId, NotificationRuleId ruleId, AlarmId alarmId) { + return notificationRequestDao.findByRuleIdAndAlarmId(tenantId, ruleId, alarmId); + } + + // ON DELETE CASCADE is used: notifications for request are deleted as well + @Override + public void deleteNotificationRequestById(TenantId tenantId, NotificationRequestId id) { + notificationRequestDao.removeById(tenantId, id.getId()); + } + + @Override + public PageData findScheduledNotificationRequests(PageLink pageLink) { + return notificationRequestDao.findAllByStatus(NotificationRequestStatus.SCHEDULED, pageLink); + } + + + private static class NotificationRequestValidator extends DataValidator { + + } + +} diff --git a/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationService.java b/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationService.java index a3ef857233..37a45a3a72 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationService.java @@ -28,6 +28,7 @@ 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.NotificationRequest; +import org.thingsboard.server.common.data.notification.NotificationRequestStatus; import org.thingsboard.server.common.data.notification.NotificationSeverity; import org.thingsboard.server.common.data.notification.NotificationStatus; import org.thingsboard.server.common.data.page.PageData; @@ -43,44 +44,8 @@ import java.util.List; @RequiredArgsConstructor public class DefaultNotificationService implements NotificationService { - private final NotificationRequestDao notificationRequestDao; private final NotificationDao notificationDao; - private final NotificationRequestValidator notificationRequestValidator = new NotificationRequestValidator(); - - @Override - public NotificationRequest saveNotificationRequest(TenantId tenantId, NotificationRequest notificationRequest) { - if (StringUtils.isBlank(notificationRequest.getNotificationReason())) { - notificationRequest.setNotificationReason(NotificationRequest.GENERAL_NOTIFICATION_REASON); - } - if (notificationRequest.getNotificationSeverity() == null) { - notificationRequest.setNotificationSeverity(NotificationSeverity.NORMAL); - } - notificationRequestValidator.validate(notificationRequest, NotificationRequest::getTenantId); - return notificationRequestDao.save(tenantId, notificationRequest); - } - - @Override - public NotificationRequest findNotificationRequestById(TenantId tenantId, NotificationRequestId id) { - return notificationRequestDao.findById(tenantId, id.getId()); - } - - @Override - public PageData findNotificationRequestsByTenantId(TenantId tenantId, PageLink pageLink) { - return notificationRequestDao.findByTenantIdAndPageLink(tenantId, pageLink); - } - - @Override - public List findNotificationRequestsByRuleIdAndAlarmId(TenantId tenantId, NotificationRuleId ruleId, AlarmId alarmId) { - return notificationRequestDao.findByRuleIdAndAlarmId(tenantId, ruleId, alarmId); - } - - // ON DELETE CASCADE is used: notifications for request are deleted as well - @Override - public void deleteNotificationRequestById(TenantId tenantId, NotificationRequestId id) { - notificationRequestDao.removeById(tenantId, id.getId()); - } - @Override public Notification saveNotification(TenantId tenantId, Notification notification) { return notificationDao.save(tenantId, notification); @@ -122,12 +87,4 @@ public class DefaultNotificationService implements NotificationService { return notificationDao.updateInfosByRequestId(tenantId, notificationRequestId, notificationInfo); } - private static class NotificationRequestValidator extends DataValidator { - - @Override - protected void validateDataImpl(TenantId tenantId, NotificationRequest notificationRequest) { - } - - } - } 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 1848e826e0..76778c64b5 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 @@ -19,6 +19,7 @@ import org.thingsboard.server.common.data.id.AlarmId; 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.NotificationRequestStatus; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.dao.Dao; @@ -31,4 +32,6 @@ public interface NotificationRequestDao extends Dao { List findByRuleIdAndAlarmId(TenantId tenantId, NotificationRuleId ruleId, AlarmId alarmId); + 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 5e1ca5d3f7..7684ed5865 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 @@ -23,6 +23,7 @@ import org.thingsboard.server.common.data.id.AlarmId; 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.NotificationRequestStatus; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.dao.DaoUtil; @@ -52,6 +53,11 @@ public class JpaNotificationRequestDao extends JpaAbstractDao findAllByStatus(NotificationRequestStatus status, PageLink pageLink) { + return DaoUtil.toPageData(notificationRequestRepository.findAllByStatus(status, DaoUtil.toPageable(pageLink))); + } + @Override protected Class getEntityClass() { return NotificationRequestEntity.class; 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 4ff69f8797..f892767a83 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.notification.NotificationRequestStatus; import org.thingsboard.server.dao.model.sql.NotificationRequestEntity; import java.util.List; @@ -37,4 +38,6 @@ public interface NotificationRequestRepository extends JpaRepository findAllByRuleIdAndAlarmId(UUID ruleId, UUID alarmId); + Page findAllByStatus(NotificationRequestStatus status, Pageable pageable); + }